From ba042781b71daffdb0efbe562b60d884d9c026d1 Mon Sep 17 00:00:00 2001 From: Prateek Gaur Date: Sat, 3 Oct 2026 22:46:31 +0000 Subject: [PATCH 1/3] Keep DELTA_BINARY_PACKED decoder state in locals The output buffer may alias decoder members of the same type, forcing repeated reloads after stores. Keep the frame and running value in locals, write back once, and preserve unsigned wrapping during reconstruction. --- cpp/src/parquet/decoder.cc | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/cpp/src/parquet/decoder.cc b/cpp/src/parquet/decoder.cc index c4d3fe5a8a5a..dc73a2eae9cb 100644 --- a/cpp/src/parquet/decoder.cc +++ b/cpp/src/parquet/decoder.cc @@ -1658,13 +1658,16 @@ class DeltaBitPackDecoder : public TypedDecoderImpl { values_decode) { ParquetException::EofException(); } + // Keep both members in locals: `buffer` may alias either one, forcing a + // reload after every output store. + UT last = static_cast(last_value_); + const UT min_delta = static_cast(min_delta_); 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(min_delta_) + static_cast(buffer[i + j]) + - static_cast(last_value_); - last_value_ = buffer[i + j]; + // Reconstruct in unsigned arithmetic so overflow wraps as specified. + last += min_delta + static_cast(buffer[i + j]); + buffer[i + j] = last; } + last_value_ = static_cast(last); } values_remaining_current_mini_block_ -= values_decode; i += values_decode; From b710c74bb9a067d55c282387d393b25f3e8dadf6 Mon Sep 17 00:00:00 2001 From: Prateek Gaur Date: Sat, 3 Oct 2026 22:47:43 +0000 Subject: [PATCH 2/3] Coalesce equal-width DELTA_BINARY_PACKED miniblocks Adjacent equal-width miniblocks form one packed stream. Decode each run with one unpack call, bounded by the current block and output request; retain the zero-width path and cover width changes, partial reads, block boundaries, and both integer widths. --- cpp/src/parquet/decoder.cc | 50 ++++++++++++++++++---- cpp/src/parquet/encoding_test.cc | 71 +++++++++++++++++++++++++++++--- 2 files changed, 106 insertions(+), 15 deletions(-) diff --git a/cpp/src/parquet/decoder.cc b/cpp/src/parquet/decoder.cc index dc73a2eae9cb..146d1c4b6173 100644 --- a/cpp/src/parquet/decoder.cc +++ b/cpp/src/parquet/decoder.cc @@ -1601,6 +1601,25 @@ class DeltaBitPackDecoder : public TypedDecoderImpl { 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(std::min(max_values, total_values_remaining_)); if (max_values == 0) { @@ -1642,35 +1661,48 @@ class DeltaBitPackDecoder : public TypedDecoderImpl { } } - int values_decode = std::min(values_remaining_current_mini_block_, - static_cast(max_values - i)); + const uint32_t values_available = static_cast(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( + 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(last_value_) + static_cast(j + 1) * static_cast(min_delta_); } - last_value_ += static_cast(values_decode) * static_cast(min_delta_); + last_value_ += + static_cast(num_values_to_decode) * static_cast(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(); } // Keep both members in locals: `buffer` may alias either one, forcing a // reload after every output store. UT last = static_cast(last_value_); const UT min_delta = static_cast(min_delta_); - for (int j = 0; j < values_decode; ++j) { + for (int j = 0; j < num_values_to_decode; ++j) { // Reconstruct in unsigned arithmetic so overflow wraps as specified. last += min_delta + static_cast(buffer[i + j]); buffer[i + j] = last; } last_value_ = static_cast(last); } - 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; diff --git a/cpp/src/parquet/encoding_test.cc b/cpp/src/parquet/encoding_test.cc index 831829e4a210..8223708641c8 100644 --- a/cpp/src/parquet/encoding_test.cc +++ b/cpp/src/parquet/encoding_test.cc @@ -23,6 +23,7 @@ #include #include #include +#include #include #include @@ -1752,7 +1753,12 @@ class TestDeltaBitPackEncoding : public TestEncodingBase { using c_type = typename Type::c_type; static constexpr int TYPE = Type::type_num; static constexpr size_t kNumRoundTrips = 3; - const std::vector kReadBatchSizes = {1, 11}; + // Keep these in sync with DeltaBitPackEncoder. + static constexpr int kValuesPerBlock = std::is_same_v ? 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 kReadBatchSizes = {1, 11, 100}; void InitBoundData(int nvalues, int repeats, c_type half_range) { num_values_ = nvalues * repeats; @@ -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 ? 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}) { @@ -2034,6 +2038,61 @@ 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& widths, T frame, int trailing_values) { + std::vector 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{1} << (width - 1)); + for (int i = 0; i < kValuesPerMiniBlock; ++i) { + current = static_cast(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(current + frame); + values.push_back(current); + } + return values; + }; + + struct Case { + const char* name; + std::vector widths; + int trailing_values; + }; + const std::vector 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(-5)}) { + ARROW_SCOPED_TRACE("case = ", c.name, ", frame = ", static_cast(frame)); + this->CheckRoundtripWithValues(make_values(c.widths, frame, c.trailing_values)); + } + } +} + // ---------------------------------------------------------------------- // Rle for Boolean encode/decode tests. From 3e46ef0d92af29f313b0ea7e0e472af2fe7c8014 Mon Sep 17 00:00:00 2001 From: Prateek Gaur Date: Sat, 3 Oct 2026 23:34:50 +0000 Subject: [PATCH 3/3] Scan DELTA_BINARY_PACKED deltas in SIMD batches Replace the dependent scalar additions with an inclusive SIMD scan for batches of at least four lanes. Retain scalar handling for narrower batches and tails, preserve unsigned wrapping, and cover every residual width, both block sizes, and wrapped totals. --- cpp/src/parquet/decoder.cc | 81 +++++++++++++++++++++++---- cpp/src/parquet/encoding_benchmark.cc | 39 +++++++++++++ cpp/src/parquet/encoding_test.cc | 36 ++++++++++++ 3 files changed, 146 insertions(+), 10 deletions(-) diff --git a/cpp/src/parquet/decoder.cc b/cpp/src/parquet/decoder.cc index 146d1c4b6173..af33faa3775f 100644 --- a/cpp/src/parquet/decoder.cc +++ b/cpp/src/parquet/decoder.cc @@ -30,6 +30,10 @@ #include #include +#if defined(ARROW_HAVE_NEON) || defined(ARROW_HAVE_SSE4_2) +# include +#endif + #include "arrow/array.h" #include "arrow/array/builder_binary.h" #include "arrow/array/builder_dict.h" @@ -1434,6 +1438,70 @@ class DictByteArrayDecoderImpl : public DictDecoderImpl { // ---------------------------------------------------------------------- // 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 +Batch InclusiveScan(Batch values) { + if constexpr (kShift < Batch::size) { + values += + xsimd::slide_left(values); + return InclusiveScan(values); + } else { + return values; + } +} +#endif + +// Reconstructs values in place using unsigned arithmetic to preserve wrapping. +template +std::make_unsigned_t ReconstructValuesFromDeltas( + T* values, int num_values, std::make_unsigned_t min_delta, + std::make_unsigned_t previous_value) { + using UT = std::make_unsigned_t; + int i = 0; + +#if defined(ARROW_HAVE_NEON) || defined(ARROW_HAVE_SSE4_2) + using Batch = xsimd::batch; + constexpr int kLanes = static_cast(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(); + 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(xsimd::batch::load_unaligned(values + i)); + batch = InclusiveScan<1>(batch + min_delta_batch) + carry; + xsimd::bitwise_cast(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(values[i]); + values[i] = static_cast(previous_value); + } + return previous_value; +} + +} // namespace + template class DeltaBitPackDecoder : public TypedDecoderImpl { public: @@ -1689,16 +1757,9 @@ class DeltaBitPackDecoder : public TypedDecoderImpl { num_values_to_decode) { ParquetException::EofException(); } - // Keep both members in locals: `buffer` may alias either one, forcing a - // reload after every output store. - UT last = static_cast(last_value_); - const UT min_delta = static_cast(min_delta_); - for (int j = 0; j < num_values_to_decode; ++j) { - // Reconstruct in unsigned arithmetic so overflow wraps as specified. - last += min_delta + static_cast(buffer[i + j]); - buffer[i + j] = last; - } - last_value_ = static_cast(last); + last_value_ = static_cast(ReconstructValuesFromDeltas( + buffer + i, num_values_to_decode, static_cast(min_delta_), + static_cast(last_value_))); } mini_block_idx_ += additional_mini_blocks; values_remaining_current_mini_block_ -= values_this_mini_block; diff --git a/cpp/src/parquet/encoding_benchmark.cc b/cpp/src/parquet/encoding_benchmark.cc index bea1a5807a2a..b50f6f69f7ed 100644 --- a/cpp/src/parquet/encoding_benchmark.cc +++ b/cpp/src/parquet/encoding_benchmark.cc @@ -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 +static auto MakeDeltaBitPackingInputIncreasing(size_t length) { + using T = typename DType::c_type; + auto numbers = std::vector(length); + ::arrow::randint(length, 0, 1000, &numbers); + T value = 0; + for (auto& number : numbers) { + value = static_cast(value + number); + number = value; + } + return numbers; +} + template static auto MakeDeltaBitPackingInputWide(size_t length) { using T = typename DType::c_type; @@ -713,6 +728,16 @@ static void BM_DeltaBitPackingEncode_Int64_Narrow(benchmark::State& state) { BM_DeltaBitPackingEncode(state, MakeDeltaBitPackingInputNarrow); } +static void BM_DeltaBitPackingEncode_Int32_Increasing(benchmark::State& state) { + BM_DeltaBitPackingEncode(state, + MakeDeltaBitPackingInputIncreasing); +} + +static void BM_DeltaBitPackingEncode_Int64_Increasing(benchmark::State& state) { + BM_DeltaBitPackingEncode(state, + MakeDeltaBitPackingInputIncreasing); +} + static void BM_DeltaBitPackingEncode_Int32_Wide(benchmark::State& state) { BM_DeltaBitPackingEncode(state, MakeDeltaBitPackingInputWide); } @@ -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); @@ -762,6 +789,16 @@ static void BM_DeltaBitPackingDecode_Int64_Narrow(benchmark::State& state) { BM_DeltaBitPackingDecode(state, MakeDeltaBitPackingInputNarrow); } +static void BM_DeltaBitPackingDecode_Int32_Increasing(benchmark::State& state) { + BM_DeltaBitPackingDecode(state, + MakeDeltaBitPackingInputIncreasing); +} + +static void BM_DeltaBitPackingDecode_Int64_Increasing(benchmark::State& state) { + BM_DeltaBitPackingDecode(state, + MakeDeltaBitPackingInputIncreasing); +} + static void BM_DeltaBitPackingDecode_Int32_Wide(benchmark::State& state) { BM_DeltaBitPackingDecode(state, MakeDeltaBitPackingInputWide); } @@ -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); diff --git a/cpp/src/parquet/encoding_test.cc b/cpp/src/parquet/encoding_test.cc index 8223708641c8..9a0f59d239d5 100644 --- a/cpp/src/parquet/encoding_test.cc +++ b/cpp/src/parquet/encoding_test.cc @@ -2093,6 +2093,42 @@ TYPED_TEST(TestDeltaBitPackEncoding, MiniblockBitWidthRuns) { } } +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; + constexpr int kBits = static_cast(sizeof(T) * 8); + constexpr int kValuesPerBlock = TestFixture::kValuesPerBlock; + + auto make_values = [](int width, T frame, int num_deltas) { + std::vector values; + values.reserve(num_deltas + 1); + const UT spread = width == kBits ? ~UT{0} : static_cast((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(current)); + for (int i = 0; i < num_deltas; ++i) { + current = static_cast(current + static_cast(frame) + + (i % 3 == 0 ? spread : UT{0})); + values.push_back(static_cast(current)); + } + return values; + }; + + for (int width = 0; width <= kBits; ++width) { + for (const T frame : {T{0}, static_cast(-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(frame), + ", num_deltas = ", num_deltas); + this->CheckRoundtripWithValues(make_values(width, frame, num_deltas)); + } + } + } +} + // ---------------------------------------------------------------------- // Rle for Boolean encode/decode tests.