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
63 changes: 49 additions & 14 deletions cpp/src/parquet/decoder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1601,6 +1601,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 +1661,48 @@ 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];
// Keep both members in locals: `buffer` may alias either one, forcing a
// reload after every output store.
UT last = static_cast<UT>(last_value_);
const UT min_delta = static_cast<UT>(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<UT>(buffer[i + j]);
buffer[i + j] = last;
}
last_value_ = static_cast<T>(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;
Expand Down
71 changes: 65 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,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<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));
}
}
}

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

Expand Down
Loading