diff --git a/cpp/src/parquet/decoder.cc b/cpp/src/parquet/decoder.cc index c4d3fe5a8a5a..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: @@ -1601,6 +1669,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 +1729,41 @@ 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]; - } + last_value_ = static_cast(ReconstructValuesFromDeltas( + buffer + i, num_values_to_decode, static_cast(min_delta_), + static_cast(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; 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 831829e4a210..9a0f59d239d5 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,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& 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)); + } + } +} + +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.