diff --git a/cpp/src/parquet/decoder.cc b/cpp/src/parquet/decoder.cc index c4d3fe5a8a5a..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,32 +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(); } - 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]; + // 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); } - 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.