From ba042781b71daffdb0efbe562b60d884d9c026d1 Mon Sep 17 00:00:00 2001 From: Prateek Gaur Date: Sat, 3 Oct 2026 22:46:31 +0000 Subject: [PATCH 1/2] 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/2] 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.