Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
126 changes: 111 additions & 15 deletions cpp/src/parquet/decoder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@
#include <utility>
#include <vector>

#if defined(ARROW_HAVE_NEON) || defined(ARROW_HAVE_SSE4_2)
# include <xsimd/xsimd.hpp>
#endif

#include "arrow/array.h"
#include "arrow/array/builder_binary.h"
#include "arrow/array/builder_dict.h"
Expand Down Expand Up @@ -1434,6 +1438,70 @@ class DictByteArrayDecoderImpl : public DictDecoderImpl<ByteArrayType> {
// ----------------------------------------------------------------------
// DELTA_BINARY_PACKED decoder

namespace {

#if defined(ARROW_HAVE_NEON) || defined(ARROW_HAVE_SSE4_2)
// Applies inclusive-scan steps recursively for power-of-two lane offsets.
template <std::size_t kShift, typename Batch>
Batch InclusiveScan(Batch values) {
if constexpr (kShift < Batch::size) {
values +=
xsimd::slide_left<kShift * sizeof(typename Batch::value_type)>(values);
return InclusiveScan<kShift * 2>(values);
} else {
return values;
}
}
#endif

// Reconstructs values in place using unsigned arithmetic to preserve wrapping.
template <typename T>
std::make_unsigned_t<T> ReconstructValuesFromDeltas(
T* values, int num_values, std::make_unsigned_t<T> min_delta,
std::make_unsigned_t<T> previous_value) {
using UT = std::make_unsigned_t<T>;
int i = 0;

#if defined(ARROW_HAVE_NEON) || defined(ARROW_HAVE_SSE4_2)
using Batch = xsimd::batch<UT>;
constexpr int kLanes = static_cast<int>(Batch::size);
// slide_left shifts toward higher lanes on the supported architectures.
// At two lanes the scan loses to the additions it replaces, so it is compiled only
// where a register holds four or more values; narrower ones use the loop below.
if constexpr (kLanes >= 4) {
// Broadcast pattern for the last lane, which carries the running value into the
// next vector without a round trip through a general-purpose register.
struct LastLane {
static constexpr unsigned get(unsigned /*index*/, unsigned size) {
return size - 1;
}
};
const auto last_lane =
xsimd::make_batch_constant<UT, LastLane, xsimd::default_arch>();
const Batch min_delta_batch(min_delta);
Batch carry(previous_value);
for (; i + kLanes <= num_values; i += kLanes) {
// Adding the frame before the scan turns its running multiple into a term the
// scan produces, rather than a multiply per lane.
Batch batch =
xsimd::bitwise_cast<UT>(xsimd::batch<T>::load_unaligned(values + i));
batch = InclusiveScan<1>(batch + min_delta_batch) + carry;
xsimd::bitwise_cast<T>(batch).store_unaligned(values + i);
carry = xsimd::swizzle(batch, last_lane);
}
previous_value = carry.get(0);
}
#endif

for (; i < num_values; ++i) {
previous_value += min_delta + static_cast<UT>(values[i]);
values[i] = static_cast<T>(previous_value);
}
return previous_value;
}

} // namespace

template <typename DType>
class DeltaBitPackDecoder : public TypedDecoderImpl<DType> {
public:
Expand Down Expand Up @@ -1601,6 +1669,25 @@ class DeltaBitPackDecoder : public TypedDecoderImpl<DType> {
values_remaining_current_mini_block_ = values_per_mini_block_;
}

// Counts additional equal-width miniblocks that can join the current unpack call.
// Miniblocks form one bitstream because each holds a multiple of 32 values.
uint32_t NumAdditionalCoalescibleMiniBlocks(uint32_t values_available) const {
// A run can only start once there is room for the rest of the current miniblock.
if (values_available < values_remaining_current_mini_block_) {
return 0;
}
const uint8_t* bit_widths = delta_bit_widths_->data();
uint32_t values_needed = values_remaining_current_mini_block_;
uint32_t count = 0;
while (mini_block_idx_ + count + 1 < mini_blocks_per_block_ &&
bit_widths[mini_block_idx_ + count + 1] == delta_bit_width_ &&
values_available - values_needed >= values_per_mini_block_) {
values_needed += values_per_mini_block_;
++count;
}
return count;
}

int GetInternal(T* buffer, int max_values) {
max_values = static_cast<int>(std::min<int64_t>(max_values, total_values_remaining_));
if (max_values == 0) {
Expand Down Expand Up @@ -1642,32 +1729,41 @@ class DeltaBitPackDecoder : public TypedDecoderImpl<DType> {
}
}

int values_decode = std::min(values_remaining_current_mini_block_,
static_cast<uint32_t>(max_values - i));
const uint32_t values_available = static_cast<uint32_t>(max_values - i);
const uint32_t values_this_mini_block =
std::min(values_remaining_current_mini_block_, values_available);
// Zero-width miniblocks bypass unpacking; coalescing them measured slower.
const uint32_t additional_mini_blocks =
delta_bit_width_ == 0
? 0
: NumAdditionalCoalescibleMiniBlocks(values_available);
// Joining another miniblock requires draining the current one.
DCHECK(additional_mini_blocks == 0 ||
values_this_mini_block == values_remaining_current_mini_block_);
const int num_values_to_decode = static_cast<int>(
values_this_mini_block + additional_mini_blocks * values_per_mini_block_);
if (delta_bit_width_ == 0) {
// Fast path that avoids a back-to-back dependency between two consecutive
// computations: we know all deltas decode to zero. We actually don't
// even need to decode them.
for (int j = 0; j < values_decode; ++j) {
for (int j = 0; j < num_values_to_decode; ++j) {
buffer[i + j] = static_cast<UT>(last_value_) +
static_cast<UT>(j + 1) * static_cast<UT>(min_delta_);
}
last_value_ += static_cast<UT>(values_decode) * static_cast<UT>(min_delta_);
last_value_ +=
static_cast<UT>(num_values_to_decode) * static_cast<UT>(min_delta_);
} else {
if (decoder_->GetBatch(delta_bit_width_, buffer + i, values_decode) !=
values_decode) {
if (decoder_->GetBatch(delta_bit_width_, buffer + i, num_values_to_decode) !=
num_values_to_decode) {
ParquetException::EofException();
}
for (int j = 0; j < values_decode; ++j) {
// Addition between min_delta, packed int and last_value should be treated as
// unsigned addition. Overflow is as expected.
buffer[i + j] = static_cast<UT>(min_delta_) + static_cast<UT>(buffer[i + j]) +
static_cast<UT>(last_value_);
last_value_ = buffer[i + j];
}
last_value_ = static_cast<T>(ReconstructValuesFromDeltas(
buffer + i, num_values_to_decode, static_cast<UT>(min_delta_),
static_cast<UT>(last_value_)));
}
values_remaining_current_mini_block_ -= values_decode;
i += values_decode;
mini_block_idx_ += additional_mini_blocks;
values_remaining_current_mini_block_ -= values_this_mini_block;
i += num_values_to_decode;
}
total_values_remaining_ -= max_values;
this->num_values_ -= max_values;
Expand Down
39 changes: 39 additions & 0 deletions cpp/src/parquet/encoding_benchmark.cc
Original file line number Diff line number Diff line change
Expand Up @@ -675,6 +675,21 @@ static auto MakeDeltaBitPackingInputNarrow(size_t length) {
return numbers;
}

// A non-decreasing random walk with deltas in [0, 1000]. Its miniblocks need 10
// bits, while Narrow's signed deltas need 11.
template <typename DType>
static auto MakeDeltaBitPackingInputIncreasing(size_t length) {
using T = typename DType::c_type;
auto numbers = std::vector<T>(length);
::arrow::randint<T, T>(length, 0, 1000, &numbers);
T value = 0;
for (auto& number : numbers) {
value = static_cast<T>(value + number);
number = value;
}
return numbers;
}

template <typename DType>
static auto MakeDeltaBitPackingInputWide(size_t length) {
using T = typename DType::c_type;
Expand Down Expand Up @@ -713,6 +728,16 @@ static void BM_DeltaBitPackingEncode_Int64_Narrow(benchmark::State& state) {
BM_DeltaBitPackingEncode<Int64Type>(state, MakeDeltaBitPackingInputNarrow<Int64Type>);
}

static void BM_DeltaBitPackingEncode_Int32_Increasing(benchmark::State& state) {
BM_DeltaBitPackingEncode<Int32Type>(state,
MakeDeltaBitPackingInputIncreasing<Int32Type>);
}

static void BM_DeltaBitPackingEncode_Int64_Increasing(benchmark::State& state) {
BM_DeltaBitPackingEncode<Int64Type>(state,
MakeDeltaBitPackingInputIncreasing<Int64Type>);
}

static void BM_DeltaBitPackingEncode_Int32_Wide(benchmark::State& state) {
BM_DeltaBitPackingEncode<Int32Type>(state, MakeDeltaBitPackingInputWide<Int32Type>);
}
Expand All @@ -725,6 +750,8 @@ BENCHMARK(BM_DeltaBitPackingEncode_Int32_Fixed)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int64_Fixed)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int32_Narrow)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int64_Narrow)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int32_Increasing)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int64_Increasing)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int32_Wide)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingEncode_Int64_Wide)->Range(MIN_RANGE, MAX_RANGE);

Expand Down Expand Up @@ -762,6 +789,16 @@ static void BM_DeltaBitPackingDecode_Int64_Narrow(benchmark::State& state) {
BM_DeltaBitPackingDecode<Int64Type>(state, MakeDeltaBitPackingInputNarrow<Int64Type>);
}

static void BM_DeltaBitPackingDecode_Int32_Increasing(benchmark::State& state) {
BM_DeltaBitPackingDecode<Int32Type>(state,
MakeDeltaBitPackingInputIncreasing<Int32Type>);
}

static void BM_DeltaBitPackingDecode_Int64_Increasing(benchmark::State& state) {
BM_DeltaBitPackingDecode<Int64Type>(state,
MakeDeltaBitPackingInputIncreasing<Int64Type>);
}

static void BM_DeltaBitPackingDecode_Int32_Wide(benchmark::State& state) {
BM_DeltaBitPackingDecode<Int32Type>(state, MakeDeltaBitPackingInputWide<Int32Type>);
}
Expand All @@ -774,6 +811,8 @@ BENCHMARK(BM_DeltaBitPackingDecode_Int32_Fixed)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int64_Fixed)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int32_Narrow)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int64_Narrow)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int32_Increasing)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int64_Increasing)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int32_Wide)->Range(MIN_RANGE, MAX_RANGE);
BENCHMARK(BM_DeltaBitPackingDecode_Int64_Wide)->Range(MIN_RANGE, MAX_RANGE);

Expand Down
107 changes: 101 additions & 6 deletions cpp/src/parquet/encoding_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
#include <functional>
#include <limits>
#include <span>
#include <type_traits>
#include <utility>
#include <vector>

Expand Down Expand Up @@ -1752,7 +1753,12 @@ class TestDeltaBitPackEncoding : public TestEncodingBase<Type> {
using c_type = typename Type::c_type;
static constexpr int TYPE = Type::type_num;
static constexpr size_t kNumRoundTrips = 3;
const std::vector<int> kReadBatchSizes = {1, 11};
// Keep these in sync with DeltaBitPackEncoder.
static constexpr int kValuesPerBlock = std::is_same_v<int32_t, c_type> ? 128 : 256;
static constexpr int kMiniBlocksPerBlock = 4;
static constexpr int kValuesPerMiniBlock = kValuesPerBlock / kMiniBlocksPerBlock;
// 100 spans several miniblocks but still ends inside one.
const std::vector<int> kReadBatchSizes = {1, 11, 100};

void InitBoundData(int nvalues, int repeats, c_type half_range) {
num_values_ = nvalues * repeats;
Expand Down Expand Up @@ -1907,11 +1913,9 @@ TYPED_TEST(TestDeltaBitPackEncoding, NonZeroPaddedMiniblockBitWidth) {
// bitwidths are actually padding bytes that may take non-conformant values
// according to the Parquet spec.

// Same values as in DeltaBitPackEncoder
constexpr int kValuesPerBlock =
std::is_same_v<int32_t, typename TypeParam::c_type> ? 128 : 256;
constexpr int kMiniBlocksPerBlock = 4;
constexpr int kValuesPerMiniBlock = kValuesPerBlock / kMiniBlocksPerBlock;
constexpr int kValuesPerBlock = TestFixture::kValuesPerBlock;
constexpr int kMiniBlocksPerBlock = TestFixture::kMiniBlocksPerBlock;
constexpr int kValuesPerMiniBlock = TestFixture::kValuesPerMiniBlock;

// num_values must be kept small enough for kHeaderLength below
for (const int num_values : {2, 62, 63, 64, 65, 95, 96, 97, 127}) {
Expand Down Expand Up @@ -2034,6 +2038,97 @@ TYPED_TEST(TestDeltaBitPackEncoding, ZeroDeltaBitWidth) {
this->CheckRoundtripWithValues(int_values);
}

TYPED_TEST(TestDeltaBitPackEncoding, MiniblockBitWidthRuns) {
// Cover equal-width runs, width changes, zero widths, block boundaries, and tails.
using T = typename TypeParam::c_type;

constexpr int kValuesPerMiniBlock = TestFixture::kValuesPerMiniBlock;

// Gives miniblock i the bit width widths[i]: alternating deltas of `frame` and
// `frame + 2^(w-1)` make w the smallest width holding the residual, and `frame` the
// smallest delta, so it is the frame the encoder stores.
auto make_values = [](const std::vector<int>& widths, T frame, int trailing_values) {
std::vector<T> values;
values.reserve(widths.size() * kValuesPerMiniBlock + trailing_values + 1);
// The first value travels in the header and contributes no delta.
T current = 0;
values.push_back(current);
for (const int width : widths) {
const T spread = width == 0 ? T{0} : static_cast<T>(T{1} << (width - 1));
for (int i = 0; i < kValuesPerMiniBlock; ++i) {
current = static_cast<T>(current + frame + (i % 2 == 0 ? T{0} : spread));
values.push_back(current);
}
}
// A tail shorter than a miniblock makes a run stop at the end of the values.
for (int i = 0; i < trailing_values; ++i) {
current = static_cast<T>(current + frame);
values.push_back(current);
}
return values;
};

struct Case {
const char* name;
std::vector<int> widths;
int trailing_values;
};
const std::vector<Case> cases = {
{"uniform widths", {4, 4, 4, 4}, 0},
{"no repeated width", {1, 8, 3, 16}, 0},
{"two runs of two", {1, 1, 8, 8}, 0},
{"run then a change", {4, 4, 4, 16}, 0},
{"zero widths first", {0, 0, 3, 3}, 0},
{"zero widths last", {3, 3, 0, 0}, 0},
{"zero width inside a run", {3, 0, 3, 3}, 0},
{"across a block boundary", {4, 4, 4, 4, 4, 4, 4, 4}, 0},
{"partial last block", {4, 4, 4, 4}, 5},
};

for (const auto& c : cases) {
for (const T frame : {T{0}, static_cast<T>(-5)}) {
ARROW_SCOPED_TRACE("case = ", c.name, ", frame = ", static_cast<int64_t>(frame));
this->CheckRoundtripWithValues(make_values(c.widths, frame, c.trailing_values));
}
}
}

TYPED_TEST(TestDeltaBitPackEncoding, ReconstructsAllWidthsAndBatchTails) {
// Cover every residual width, SIMD batch tail, block boundary, and unsigned wrap.
using T = typename TypeParam::c_type;
using UT = std::make_unsigned_t<T>;
constexpr int kBits = static_cast<int>(sizeof(T) * 8);
constexpr int kValuesPerBlock = TestFixture::kValuesPerBlock;

auto make_values = [](int width, T frame, int num_deltas) {
std::vector<T> values;
values.reserve(num_deltas + 1);
const UT spread = width == kBits ? ~UT{0} : static_cast<UT>((UT{1} << width) - 1);
// Two deltas in three sit at the frame, so it is the smallest in every miniblock.
UT current = 0;
values.push_back(static_cast<T>(current));
for (int i = 0; i < num_deltas; ++i) {
current = static_cast<UT>(current + static_cast<UT>(frame) +
(i % 3 == 0 ? spread : UT{0}));
values.push_back(static_cast<T>(current));
}
return values;
};

for (int width = 0; width <= kBits; ++width) {
for (const T frame : {T{0}, static_cast<T>(-5), T{7}}) {
// 16-23 leaves every remainder for group sizes up to eight. The last length
// crosses a block boundary with a remainder left, at either block size.
for (const int num_deltas :
{16, 17, 18, 19, 20, 21, 22, 23, kValuesPerBlock + 73}) {
ARROW_SCOPED_TRACE("width = ", width, ", frame = ", static_cast<int64_t>(frame),
", num_deltas = ", num_deltas);
this->CheckRoundtripWithValues(make_values(width, frame, num_deltas));
}
}
}
}

// ----------------------------------------------------------------------
// Rle for Boolean encode/decode tests.

Expand Down
Loading