Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
62 commits
Select commit Hold shift + click to select a range
e17eec9
Move TypedColReaderImpl::Skip
AntoinePrv Jul 9, 2026
feed609
Break TypedRecordReader inheritance on TypedColumnReader
AntoinePrv Jul 9, 2026
202b354
Add ReadValuesBuffer
AntoinePrv Jul 10, 2026
7829db2
Split ColumnReaderImplBase definition and declaration
AntoinePrv Jul 10, 2026
63c56a8
Add DefRepLevelsDecoder
AntoinePrv Jul 10, 2026
46b395f
Explicit int32_t for LevelDecoder
AntoinePrv Jul 10, 2026
f46f89a
Change ValuesBuffer API
AntoinePrv Jul 15, 2026
2928a80
Add StackedLevelDecoder and prepare for modularity
AntoinePrv Jul 15, 2026
7fe2381
Remove DefRepLevelsDecoder
AntoinePrv Jul 15, 2026
17ac56b
Split TypedRecordReader implementation
AntoinePrv Jul 16, 2026
7ce96cb
First RequiredTypedRecordReader
AntoinePrv Jul 16, 2026
b5f0424
Remove RecordReader protected members
AntoinePrv Jul 16, 2026
bb51e19
Fully use RequiredRecordReader
AntoinePrv Jul 16, 2026
79eb654
Rename ReadValues -> ValueSink
AntoinePrv Jul 17, 2026
22372da
Split ReserveValuesAndIsValid
AntoinePrv Jul 17, 2026
f9dec99
Devirtualize hooks with CRTP
AntoinePrv Jul 20, 2026
1d922a6
Simplify template
AntoinePrv Jul 20, 2026
463ad5d
Remove StackedLevelDecoder
AntoinePrv Jul 20, 2026
6d86e26
Remove mark_as_written
AntoinePrv Jul 20, 2026
4823172
Make some of ColumnReaderImplBase private
AntoinePrv Jul 20, 2026
931d075
Fix horrible dictionnary handling
AntoinePrv Jul 20, 2026
3320b2b
Simplify and rename ColumnChunkReader
AntoinePrv Jul 20, 2026
392b71a
Add RleBitPackedToBitmapDecoder::CountUpTo
AntoinePrv Jul 21, 2026
343e0bf
Add BitPackedToBitmapDecoder
AntoinePrv Jul 21, 2026
0a2e3c1
Add RleBitPackedToBitmapDecoder::Advance
AntoinePrv Jul 21, 2026
af37312
Add LevelToBitmapDecoder
AntoinePrv Jul 21, 2026
fadc841
Add ValiditySinkBuffer
AntoinePrv Jul 22, 2026
2bf286f
Auto reserve space in sinks:
AntoinePrv Jul 22, 2026
f24efa1
Fix integer types
AntoinePrv Jul 22, 2026
8043b45
Nit in ColumnChunkReader
AntoinePrv Jul 22, 2026
2444c47
Eagerly read dict
AntoinePrv Jul 23, 2026
a3cb538
Cleanups
AntoinePrv Jul 24, 2026
c297eed
Add FlatOptionalTypedRecordReader
AntoinePrv Jul 24, 2026
778c979
Add some doc
AntoinePrv Jul 27, 2026
ad6e524
Factorize ReadDictionary
AntoinePrv Jul 29, 2026
f4ae4fc
Simplify read_dense_for_nullable
AntoinePrv Jul 29, 2026
fe8bf3d
Factor Sinks with DataSinkBuffer
AntoinePrv Jul 30, 2026
cf7aa77
Use LevelBufferSink in TypedRecordReader
AntoinePrv Jul 30, 2026
056affa
Relax record reader level tests flat optional
AntoinePrv Jul 30, 2026
7e251a3
Fix compiler errrors
AntoinePrv Jul 31, 2026
3e6eeb9
Safe decoded number check and func rename
AntoinePrv Jul 31, 2026
f268787
Explicit values_written contract
AntoinePrv Jul 31, 2026
dbeaec0
Fix LevelDecoder::max_level type
AntoinePrv Jul 31, 2026
f7c56ce
Fix integer convertion
AntoinePrv Jul 31, 2026
cca5726
Factorize level decoders
AntoinePrv Aug 3, 2026
c8b1c2a
Fixt typos
AntoinePrv Aug 6, 2026
6e4465f
Move level decoder utility and rename file
AntoinePrv Sep 29, 2026
8089237
Move level decoder adapters to level_decoder.h
AntoinePrv Sep 29, 2026
47e0874
Fix small typos
AntoinePrv Sep 29, 2026
baba303
Fix issues in level counting
AntoinePrv Sep 29, 2026
31d2007
Fix extra alloc
AntoinePrv Oct 1, 2026
9e3f8bc
Add RecordReader::Make flag
AntoinePrv Oct 6, 2026
07ee3c5
Throw on FlatOptional def_levels
AntoinePrv Oct 6, 2026
02c9b2c
Always call RecordReader with flat optiona optimization
AntoinePrv Oct 6, 2026
d833502
Addition flat optional test
AntoinePrv Oct 6, 2026
fda90cb
Detemplate kdic
AntoinePrv Oct 6, 2026
dbb7438
Add flat optiona bse in flba and byte array
AntoinePrv Oct 7, 2026
1a205c8
Ad more flba and array test
AntoinePrv Oct 7, 2026
5516651
Make RecordReader::Make params struct public
AntoinePrv Oct 8, 2026
d911101
Run RecordReader tests with and without optimizations
AntoinePrv Oct 8, 2026
52e062c
Rename flat_optiona_optimization as allow_levels_optimization
AntoinePrv Oct 8, 2026
74b3576
Add level optim option to RowGroupReader
AntoinePrv Oct 8, 2026
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
157 changes: 138 additions & 19 deletions cpp/src/arrow/util/rle_bitmap_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

#include <algorithm>
#include <cstring>
#include <limits>

#include "arrow/util/bit_util.h"
#include "arrow/util/bitmap_ops.h"
Expand Down Expand Up @@ -112,6 +113,18 @@ class RleRunToBitmapDecoder {
return n_vals;
}

/// Advance and count the number of occurrences of `value`.
///
/// The count is limited to at most the next `batch_size` items.
/// @return The matching value count and number of elements that were processed.
RleCountUpToResult CountUpTo(bool value, rle_size_t batch_size) {
const auto n_vals = Advance(batch_size);
return {
.matching_count = n_vals * (value == this->value()),
.processed_count = n_vals,
};
}

/// Read the next value into `out` and return false if there are no more.
[[nodiscard]] bool Get(BitmapSpanMut out) { return GetBatch(out, 1) == 1; }

Expand Down Expand Up @@ -229,6 +242,19 @@ class BitPackedRunToBitmapDecoder {
return n_vals;
}

/// Advance and count the number of occurrences of `value`.
///
/// The count is limited to at most the next `batch_size` items.
/// @return The matching value count and number of elements that were processed.
RleCountUpToResult CountUpTo(bool value, rle_size_t batch_size) {
const auto n_vals = GetAdvanceCapacity(batch_size);
const auto set_count = static_cast<rle_size_t>(arrow::internal::CountSetBits(
unread_values_ptr(), unread_values_bit_offset(), n_vals));
const auto matching_count = value ? set_count : n_vals - set_count;
ARROW_UNUSED(Advance(n_vals));
return {.matching_count = matching_count, .processed_count = n_vals};
}

/// Get the next value and return false if there are no more.
[[nodiscard]] bool Get(BitmapSpanMut out) { return GetBatch(out, 1) == 1; }

Expand Down Expand Up @@ -282,11 +308,23 @@ class RleBitPackedToBitmapDecoder {
/// Whether there is still runs to iterate over.
bool exhausted() const { return (run_remaining() == 0) && parser_.exhausted(); }

/// Advance by as many values as provided or until exhaustion of the decoder.
/// Return the number of values skipped.
[[nodiscard]] rle_size_t Advance(rle_size_t batch_size);

/// Get a batch of values return the number of decoded elements.
/// May write fewer elements to the output than requested if there are not enough
/// values left or if an error occurred.
[[nodiscard]] rle_size_t GetBatch(BitmapSpanMut out, rle_size_t batch_size);

/// Advance and count the number of occurrences of `value`.
///
/// The count is limited to at most the next `batch_size` items. Fewer than
/// `batch_size` elements may be processed if there are not enough values left
/// or if an error occurred.
/// @return The matching value count and number of elements that were processed.
RleCountUpToResult CountUpTo(bool value, rle_size_t batch_size);

private:
RleBitPackedParser parser_ = {};
std::variant<RleRunToBitmapDecoder, BitPackedRunToBitmapDecoder> decoder_ = {};
Expand All @@ -296,10 +334,60 @@ class RleBitPackedToBitmapDecoder {
return std::visit([](const auto& dec) { return dec.remaining(); }, decoder_);
}

/// Get a batch of values from the current run and return the number elements read.
[[nodiscard]] rle_size_t RunGetBatch(BitmapSpanMut out, rle_size_t batch_size) {
return std::visit([&](auto& dec) { return dec.GetBatch(out, batch_size); }, decoder_);
/// Process data in the current run and subsequent ones.
///
/// `func` is called as `func(decoder, run_batch_size)` where `decoder` is a
/// statically-typed run decoder (not the variant).
/// Must return the number of values it processed.
///
/// Return the number of values processed.
template <typename Callable>
rle_size_t ProcessValues(Callable&& func, rle_size_t batch_size);
};

/// Minimal decoder for legacy bit packed encoding (BIT_PACKED = 4) into a bitmap.
///
/// This is the bitmap counterpart of ``BitPackedDecoder``: it decodes a single
/// bit-packed run of booleans (each value stored on a single bit) directly into an
/// Arrow validity bitmap, without the RLE framing.
///
/// The number of values that the decoder can represent is up to 2^31 - 1.
class BitPackedToBitmapDecoder : private BitPackedRunToBitmapDecoder {
private:
using Base = BitPackedRunToBitmapDecoder;

public:
using Base::Advance;
using Base::CountUpTo;
using Base::Get;
using Base::GetBatch;
using Base::remaining;

BitPackedToBitmapDecoder() noexcept = default;

/// Create a decoder object.
///
/// @param data and data_size are the raw bytes to decode.
/// @param value_count is the number of values in the run.
BitPackedToBitmapDecoder(const uint8_t* data, rle_size_t data_size,
rle_size_t value_count) noexcept {
Reset(data, data_size, value_count);
}

void Reset(const uint8_t* data, rle_size_t data_size, rle_size_t value_count) noexcept {
ARROW_DCHECK_GE(value_count, 0);
ARROW_DCHECK_LE(value_count, std::numeric_limits<rle_size_t>::max());
const auto run = BitPackedRun{
/* data= */ data,
/* value_count= */ value_count,
/* value_bit_width= */ 1,
/* max_read_bytes= */ data_size,
};
return Base::Reset(run);
}

/// Whether there is still values to iterate over.
bool exhausted() const { return Base::remaining() == 0; }
};

/************************************************
Expand All @@ -320,50 +408,81 @@ struct RleBitPackedToBitmapDecoderGetDecoder<BitPackedRun> {
using type = BitPackedRunToBitmapDecoder;
};

inline auto RleBitPackedToBitmapDecoder::GetBatch(BitmapSpanMut out,
rle_size_t batch_size) -> rle_size_t {
template <typename Callable>
auto RleBitPackedToBitmapDecoder::ProcessValues(Callable&& func,
rle_size_t batch_size) -> rle_size_t {
using ControlFlow = RleBitPackedParser::ControlFlow;

if (ARROW_PREDICT_FALSE(batch_size == 0 || exhausted())) {
return 0;
}

rle_size_t values_read = 0;
rle_size_t values_processed = 0;

// Remaining from a previous call that would have left some unread data from a run.
if (ARROW_PREDICT_FALSE(run_remaining() > 0)) {
const auto read = RunGetBatch(out, batch_size);
values_read += read;
const auto processed =
std::visit([&](auto& decoder) { return func(decoder, batch_size); }, decoder_);
values_processed += processed;

// Either we fulfilled all the batch to be read or we finished remaining run.
if (ARROW_PREDICT_FALSE(values_read == batch_size)) {
return values_read;
if (ARROW_PREDICT_FALSE(values_processed == batch_size)) {
return values_processed;
}
ARROW_DCHECK(run_remaining() == 0);
}

parser_.ParseWithCallable([&](auto run) {
using RunDecoder = RleBitPackedToBitmapDecoderGetDecoder<decltype(run)>::type;

ARROW_DCHECK_LT(values_read, batch_size);
ARROW_DCHECK_LT(values_processed, batch_size);
RunDecoder decoder(run);
// The output span carries its own bit offset, so advancing it past the values
// already written keeps successive runs correctly aligned in the bitmap.
const auto read =
decoder.GetBatch(out.NewStartingAt(values_read), batch_size - values_read);
ARROW_DCHECK_LE(read, batch_size - values_read);
values_read += read;
const auto read = func(decoder, batch_size - values_processed);
ARROW_DCHECK_LE(read, batch_size - values_processed);
values_processed += read;

// Stop reading and store remaining decoder
if (ARROW_PREDICT_FALSE(values_read == batch_size || read == 0)) {
if (ARROW_PREDICT_FALSE(values_processed == batch_size || read == 0)) {
decoder_ = std::move(decoder);
return ControlFlow::Break;
}

return ControlFlow::Continue;
});

return values_read;
return values_processed;
}

inline auto RleBitPackedToBitmapDecoder::Advance(rle_size_t batch_size) -> rle_size_t {
return ProcessValues(
[](auto& decoder, rle_size_t run_batch_size) {
return decoder.Advance(run_batch_size);
},
batch_size);
}

inline auto RleBitPackedToBitmapDecoder::GetBatch(BitmapSpanMut out,
rle_size_t batch_size) -> rle_size_t {
return ProcessValues(
[&out](auto& decoder, rle_size_t run_batch_size) {
const auto read = decoder.GetBatch(out, run_batch_size);
out = out.NewStartingAt(read);
return read;
},
batch_size);
}

inline auto RleBitPackedToBitmapDecoder::CountUpTo(bool value, rle_size_t batch_size)
-> RleCountUpToResult {
rle_size_t matching_count = 0;
const rle_size_t processed_count = ProcessValues(
[value, &matching_count](auto& decoder, rle_size_t run_batch_size) {
const auto result = decoder.CountUpTo(value, run_batch_size);
matching_count += result.matching_count;
return result.processed_count;
},
batch_size);
return {.matching_count = matching_count, .processed_count = processed_count};
}

} // namespace arrow::util
Loading
Loading