This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git
The following commit(s) were added to refs/heads/main by this push:
new 75544f7 feat(parquet): skip page headers of unselected pages via
OffsetIndex-based direct read plan (#167)
75544f7 is described below
commit 75544f79df885fe76ff3309ce530626ce18cb35c
Author: Zhou Hongfeng <[email protected]>
AuthorDate: Mon Aug 3 13:05:53 2026 +0800
feat(parquet): skip page headers of unselected pages via OffsetIndex-based
direct read plan (#167)
* fix(parquet): skip reading page header for unneeded pages
* style: clang tidy
* fix
---
cmake_modules/arrow.diff | 209 ++++++++
.../parquet/page_filtered_row_group_reader.cpp | 206 ++++++--
.../parquet/page_filtered_row_group_reader.h | 20 +-
.../page_filtered_row_group_reader_test.cpp | 577 ++++++++++++++++++++-
4 files changed, 940 insertions(+), 72 deletions(-)
diff --git a/cmake_modules/arrow.diff b/cmake_modules/arrow.diff
index c8aa0dd..f86f36e 100644
--- a/cmake_modules/arrow.diff
+++ b/cmake_modules/arrow.diff
@@ -741,3 +741,212 @@ diff --git a/cpp/src/arrow/io/interfaces.h
b/cpp/src/arrow/io/interfaces.h
Status TransferColumnData(::parquet::internal::RecordReader* reader,
const std::shared_ptr<::arrow::Field>& value_field,
const ColumnDescriptor* descr, ::arrow::MemoryPool*
pool,
+diff --git a/cpp/src/parquet/column_reader.h b/cpp/src/parquet/column_reader.h
+--- a/cpp/src/parquet/column_reader.h
++++ b/cpp/src/parquet/column_reader.h
+@@ -76,6 +76,18 @@ struct PARQUET_EXPORT DataPageStats {
+ std::optional<int32_t> num_rows;
+ };
+
++/// \brief Identifies a data page that PageReader should read directly.
++///
++/// The offset is relative to the beginning of the column chunk stream passed
to
++/// PageReader::Open. The compressed size includes both the serialized page
header
++/// and the compressed page body. The ordinal is the original data page
ordinal in
++/// the column chunk and does not include the dictionary page.
++struct PARQUET_EXPORT DataPageReadPlanEntry {
++ int32_t page_ordinal;
++ int64_t offset;
++ int32_t compressed_page_size;
++};
++
+ class PARQUET_EXPORT LevelDecoder {
+ public:
+ LevelDecoder();
+@@ -147,9 +159,21 @@ class PARQUET_EXPORT PageReader {
+ // ApplicationVersion::HasCorrectStatistics().
+ // \note API EXPERIMENTAL
+ void set_data_page_filter(DataPageFilter data_page_filter) {
++ if (data_page_read_plan_enabled_) {
++ throw ParquetException(
++ "data_page_filter and data_page_read_plan cannot be enabled
together");
++ }
+ data_page_filter_ = std::move(data_page_filter);
+ }
+
++ /// Configure PageReader to jump directly to selected data pages before
reading
++ /// their headers. `first_data_page_offset` and each entry offset are
relative to
++ /// the beginning of the column chunk stream. Dictionary pages before
++ /// `first_data_page_offset` are still read normally.
++ // \note API EXPERIMENTAL
++ void set_data_page_read_plan(int64_t first_data_page_offset,
++ std::vector<DataPageReadPlanEntry> data_pages);
++
+ // @returns: shared_ptr<Page>(nullptr) on EOS, std::shared_ptr<Page>
+ // containing new Page otherwise
+ //
+@@ -162,6 +186,11 @@ class PARQUET_EXPORT PageReader {
+ protected:
+ // Callback that decides if we should skip a page or not.
+ DataPageFilter data_page_filter_;
++
++ bool data_page_read_plan_enabled_ = false;
++ int64_t first_data_page_offset_ = 0;
++ std::vector<DataPageReadPlanEntry> data_page_read_plan_;
++ size_t next_data_page_ = 0;
+ };
+
+ class PARQUET_EXPORT ColumnReader {
+diff --git a/cpp/src/parquet/column_reader.cc
b/cpp/src/parquet/column_reader.cc
+--- a/cpp/src/parquet/column_reader.cc
++++ b/cpp/src/parquet/column_reader.cc
+@@ -207,6 +207,39 @@ ReaderProperties default_reader_properties() {
+ return default_reader_properties;
+ }
+
++void PageReader::set_data_page_read_plan(
++ int64_t first_data_page_offset,
++ std::vector<DataPageReadPlanEntry> data_pages) {
++ if (data_page_filter_) {
++ throw ParquetException(
++ "data_page_filter and data_page_read_plan cannot be enabled
together");
++ }
++ if (first_data_page_offset < 0) {
++ throw ParquetException("Invalid negative first data page offset");
++ }
++
++ int64_t previous_end = first_data_page_offset;
++ int32_t previous_ordinal = -1;
++ for (const auto& page : data_pages) {
++ int64_t page_end;
++ if (page.page_ordinal < 0 || page.offset < first_data_page_offset ||
++ page.compressed_page_size <= 0 ||
++ AddWithOverflow(page.offset, page.compressed_page_size, &page_end)) {
++ throw ParquetException("Invalid data page read plan entry");
++ }
++ if (page.offset < previous_end || page.page_ordinal <= previous_ordinal) {
++ throw ParquetException("Data page read plan entries must be ordered");
++ }
++ previous_end = page_end;
++ previous_ordinal = page.page_ordinal;
++ }
++
++ data_page_read_plan_enabled_ = true;
++ first_data_page_offset_ = first_data_page_offset;
++ data_page_read_plan_ = std::move(data_pages);
++ next_data_page_ = 0;
++}
++
+ namespace {
+
+ // Extracts encoded statistics from V1 and V2 data page headers
+@@ -430,9 +463,43 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+
+ // Loop here because there may be unhandled page types that we skip until
+ // finding a page that we do know what to do with
+- while (seen_num_values_ < total_num_values_) {
++ while (data_page_read_plan_enabled_ || seen_num_values_ <
total_num_values_) {
++ const DataPageReadPlanEntry* planned_data_page = nullptr;
++ uint32_t page_header_limit = max_page_header_size_;
++
++ if (data_page_read_plan_enabled_) {
++ if (next_data_page_ >= data_page_read_plan_.size()) {
++ return nullptr;
++ }
++
++ PARQUET_ASSIGN_OR_THROW(int64_t current_position, stream_->Tell());
++ if (current_position < first_data_page_offset_) {
++ page_header_limit = static_cast<uint32_t>(std::min<int64_t>(
++ page_header_limit, first_data_page_offset_ - current_position));
++ } else {
++ planned_data_page = &data_page_read_plan_[next_data_page_];
++ if (current_position > planned_data_page->offset) {
++ throw ParquetException("Data page read plan points behind stream
position");
++ }
++ PARQUET_THROW_NOT_OK(
++ stream_->Advance(planned_data_page->offset - current_position));
++ PARQUET_ASSIGN_OR_THROW(int64_t target_position, stream_->Tell());
++ if (target_position != planned_data_page->offset) {
++ throw ParquetException("Failed to seek to planned data page");
++ }
++ page_ordinal_ = planned_data_page->page_ordinal;
++ page_header_limit = static_cast<uint32_t>(std::min<int64_t>(
++ page_header_limit, planned_data_page->compressed_page_size));
++ }
++ }
++
++ if (page_header_limit == 0) {
++ throw ParquetException("No bytes available for page header");
++ }
++
+ uint32_t header_size = 0;
+- uint32_t allowed_page_size = kDefaultPageHeaderSize;
++ uint32_t allowed_page_size =
++ std::min<uint32_t>(kDefaultPageHeaderSize, page_header_limit);
+
+ // Page headers can be very large because of page statistics
+ // We try to deserialize a larger buffer progressively
+@@ -458,11 +525,12 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+ // Failed to deserialize. Double the allowed page header size and try
again
+ std::stringstream ss;
+ ss << e.what();
+- allowed_page_size *= 2;
+- if (allowed_page_size > max_page_header_size_) {
++ if (allowed_page_size >= page_header_limit) {
+ ss << "Deserializing page header failed.\n";
+ throw ParquetException(ss.str());
+ }
++ allowed_page_size =
++ std::min<uint32_t>(allowed_page_size * 2, page_header_limit);
+ }
+ }
+ // Advance the stream offset
+@@ -474,6 +542,20 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+ throw ParquetException("Invalid page header");
+ }
+
++ const PageType::type page_type = LoadEnumSafe(¤t_page_header_.type);
++ if (planned_data_page != nullptr) {
++ if (page_type != PageType::DATA_PAGE && page_type !=
PageType::DATA_PAGE_V2) {
++ throw ParquetException("Data page read plan points to a non-data
page");
++ }
++ int64_t total_compressed_size;
++ if (AddWithOverflow(static_cast<int64_t>(header_size),
++ static_cast<int64_t>(compressed_len),
++ &total_compressed_size) ||
++ total_compressed_size != planned_data_page->compressed_page_size) {
++ throw ParquetException("Planned data page size does not match page
header");
++ }
++ }
++
+ EncodedStatistics data_page_statistics;
+ if (ShouldSkipPage(&data_page_statistics)) {
+ PARQUET_THROW_NOT_OK(stream_->Advance(compressed_len));
+@@ -494,8 +576,6 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+ ParquetException::EofException(ss.str());
+ }
+
+- const PageType::type page_type = LoadEnumSafe(¤t_page_header_.type);
+-
+ if (properties_.page_checksum_verification() &&
current_page_header_.__isset.crc &&
+ PageCanUseChecksum(page_type)) {
+ // verify crc
+@@ -534,6 +614,9 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+
LoadEnumSafe(&dict_header.encoding),
+ is_sorted);
+ } else if (page_type == PageType::DATA_PAGE) {
++ if (planned_data_page != nullptr) {
++ ++next_data_page_;
++ }
+ ++page_ordinal_;
+ const format::DataPageHeader& header =
current_page_header_.data_page_header;
+ page_buffer =
+@@ -545,6 +628,9 @@ std::shared_ptr<Page> SerializedPageReader::NextPage() {
+ LoadEnumSafe(&header.repetition_level_encoding), uncompressed_len,
+ std::move(data_page_statistics));
+ } else if (page_type == PageType::DATA_PAGE_V2) {
++ if (planned_data_page != nullptr) {
++ ++next_data_page_;
++ }
+ ++page_ordinal_;
+ const format::DataPageHeaderV2& header =
current_page_header_.data_page_header_v2;
diff --git a/src/paimon/format/parquet/page_filtered_row_group_reader.cpp
b/src/paimon/format/parquet/page_filtered_row_group_reader.cpp
index 6b8f9a2..c8129db 100644
--- a/src/paimon/format/parquet/page_filtered_row_group_reader.cpp
+++ b/src/paimon/format/parquet/page_filtered_row_group_reader.cpp
@@ -20,6 +20,8 @@
#include "paimon/format/parquet/page_filtered_row_group_reader.h"
#include <algorithm>
+#include <limits>
+#include <optional>
#include "arrow/array.h"
#include "arrow/builder.h"
@@ -41,6 +43,56 @@ namespace paimon::parquet {
namespace {
+struct DataPageLayout {
+ int64_t column_chunk_offset;
+ int64_t first_data_page_offset;
+};
+
+int64_t GetColumnChunkOffset(const ::parquet::ColumnChunkMetaData&
column_chunk) {
+ int64_t column_chunk_offset = column_chunk.data_page_offset();
+ if (column_chunk.has_dictionary_page() &&
column_chunk.dictionary_page_offset() > 0 &&
+ column_chunk.dictionary_page_offset() < column_chunk_offset) {
+ column_chunk_offset = column_chunk.dictionary_page_offset();
+ }
+ return column_chunk_offset;
+}
+
+std::optional<DataPageLayout> GetDataPageLayout(
+ const ::parquet::ColumnChunkMetaData& column_chunk,
+ const std::shared_ptr<::parquet::OffsetIndex>& offset_index, int64_t
row_group_row_count) {
+ const auto& page_locations = offset_index->page_locations();
+ if (page_locations.empty() || row_group_row_count <= 0 ||
+ page_locations.front().first_row_index != 0) {
+ return std::nullopt;
+ }
+
+ const int64_t column_chunk_offset = GetColumnChunkOffset(column_chunk);
+ const int64_t first_data_page_offset = page_locations.front().offset;
+ const int64_t column_chunk_size = column_chunk.total_compressed_size();
+ if (column_chunk_offset < 0 || column_chunk_size <= 0 ||
+ column_chunk_offset > std::numeric_limits<int64_t>::max() -
column_chunk_size ||
+ first_data_page_offset < column_chunk_offset) {
+ return std::nullopt;
+ }
+ const int64_t column_chunk_end = column_chunk_offset + column_chunk_size;
+
+ int64_t previous_page_end = first_data_page_offset;
+ int64_t previous_first_row = -1;
+ for (const auto& page : page_locations) {
+ if (page.offset < previous_page_end || page.compressed_page_size <= 0
||
+ page.offset > std::numeric_limits<int64_t>::max() -
page.compressed_page_size ||
+ page.offset + page.compressed_page_size > column_chunk_end ||
+ page.first_row_index <= previous_first_row || page.first_row_index
< 0 ||
+ page.first_row_index >= row_group_row_count) {
+ return std::nullopt;
+ }
+ previous_page_end = page.offset + page.compressed_page_size;
+ previous_first_row = page.first_row_index;
+ }
+
+ return DataPageLayout{column_chunk_offset, first_data_page_offset};
+}
+
/// Wraps an arrow::Table + TableBatchReader as a RecordBatchReader so the
caller can
/// stream batches while ensuring every returned array offset is zero. The
Table is held
/// to keep its ChunkedArrays alive for the inner TableBatchReader.
@@ -79,27 +131,32 @@ class TableRecordBatchReader : public
arrow::RecordBatchReader {
std::shared_ptr<arrow::MemoryPool> pool_;
};
-/// A FileColumnIterator that installs a data_page_filter on every PageReader
it
-/// produces, enabling I/O-level page skipping. The base class handles row
group
-/// iteration; this subclass only decorates the PageReader returned by
NextChunk().
+/// A FileColumnIterator that installs a direct data page read plan on every
PageReader
+/// it produces. The base class handles row group iteration; this subclass only
+/// decorates the PageReader returned by NextChunk().
class PageFilteringColumnIterator : public
::parquet::arrow::FileColumnIterator {
public:
- PageFilteringColumnIterator(
- int column_index, ::parquet::ParquetFileReader* reader,
std::vector<int> row_groups,
- std::function<bool(const ::parquet::DataPageStats&)> data_page_filter)
+ PageFilteringColumnIterator(int column_index,
::parquet::ParquetFileReader* reader,
+ std::vector<int> row_groups, bool
has_data_page_read_plan,
+ int64_t first_data_page_offset,
+ std::vector<::parquet::DataPageReadPlanEntry>
data_pages)
: FileColumnIterator(column_index, reader, std::move(row_groups)),
- data_page_filter_(std::move(data_page_filter)) {}
+ has_data_page_read_plan_(has_data_page_read_plan),
+ first_data_page_offset_(first_data_page_offset),
+ data_pages_(std::move(data_pages)) {}
std::unique_ptr<::parquet::PageReader> NextChunk() override {
std::unique_ptr<::parquet::PageReader> page_reader =
FileColumnIterator::NextChunk();
- if (page_reader && data_page_filter_) {
- page_reader->set_data_page_filter(data_page_filter_);
+ if (page_reader && has_data_page_read_plan_) {
+ page_reader->set_data_page_read_plan(first_data_page_offset_,
data_pages_);
}
return page_reader;
}
private:
- std::function<bool(const ::parquet::DataPageStats&)> data_page_filter_;
+ bool has_data_page_read_plan_;
+ int64_t first_data_page_offset_;
+ std::vector<::parquet::DataPageReadPlanEntry> data_pages_;
};
} // namespace
@@ -114,22 +171,32 @@ std::pair<int64_t, int64_t>
PageFilteredRowGroupReader::GetPageRowRange(
return {first_row, last_row};
}
-std::function<bool(const ::parquet::DataPageStats&)>
PageFilteredRowGroupReader::MakePageFilter(
+std::optional<PageFilteredRowGroupReader::DataPageReadPlan>
+PageFilteredRowGroupReader::MakeDataPageReadPlan(
const RowRanges& row_ranges, const
std::shared_ptr<::parquet::OffsetIndex>& offset_index,
- int64_t row_group_row_count) {
- auto page_counter = std::make_shared<int32_t>(0);
+ const ::parquet::ColumnChunkMetaData& column_chunk, int64_t
row_group_row_count) {
+ std::optional<DataPageLayout> layout =
+ GetDataPageLayout(column_chunk, offset_index, row_group_row_count);
+ if (!layout) {
+ return std::nullopt;
+ }
+
const auto& page_locations = offset_index->page_locations();
auto num_pages = static_cast<int32_t>(page_locations.size());
+ std::vector<::parquet::DataPageReadPlanEntry> data_pages;
+ data_pages.reserve(page_locations.size());
- return [row_ranges, page_locations, num_pages, row_group_row_count,
- page_counter](const ::parquet::DataPageStats& /*stats*/) -> bool {
- int32_t page_idx = (*page_counter)++;
- if (page_idx >= num_pages) {
- return false;
- }
+ for (int32_t page_idx = 0; page_idx < num_pages; ++page_idx) {
auto [first_row, last_row] = GetPageRowRange(page_locations, page_idx,
row_group_row_count);
- return !row_ranges.IsOverlapping(first_row, last_row);
- };
+ if (row_ranges.IsOverlapping(first_row, last_row)) {
+ const auto& page = page_locations[page_idx];
+ data_pages.push_back(
+ {page_idx, page.offset - layout->column_chunk_offset,
page.compressed_page_size});
+ }
+ }
+
+ return DataPageReadPlan{layout->first_data_page_offset -
layout->column_chunk_offset,
+ std::move(data_pages)};
}
std::pair<RowRanges, int64_t>
PageFilteredRowGroupReader::ComputeCompressedRowRanges(
@@ -147,7 +214,7 @@ std::pair<RowRanges, int64_t>
PageFilteredRowGroupReader::ComputeCompressedRowRa
int64_t page_size = page_to - page_from + 1;
if (!original_ranges.IsOverlapping(page_from, page_to)) {
- // Page will be skipped by data_page_filter, not in compressed
space
+ // Page will be skipped by the direct read plan, not in compressed
space
continue;
}
@@ -220,21 +287,40 @@ Result<std::shared_ptr<arrow::ChunkedArray>>
PageFilteredRowGroupReader::ReadFil
int32_t row_group_index, int32_t field_index, const std::vector<int32_t>&
column_indices,
const RowRanges& row_ranges, int64_t row_group_row_count,
::parquet::arrow::FileReader* arrow_file_reader) {
- // Factory: set data_page_filter on every leaf (per-leaf OffsetIndex).
- // data_page_filter enables I/O-level page skipping for all leaves.
+ // Factory: set a direct data page read plan on every leaf (per-leaf
OffsetIndex).
+ // The plan lets Arrow jump over unselected page headers as well as page
bodies.
auto factory =
[row_group_index, &rg_page_index_reader, &row_ranges,
row_group_row_count](
int col_idx,
::parquet::ParquetFileReader* reader) ->
::parquet::arrow::FileColumnIterator* {
- std::function<bool(const ::parquet::DataPageStats&)> data_page_filter;
+ bool has_data_page_read_plan = false;
+ int64_t first_data_page_offset = 0;
+ std::vector<::parquet::DataPageReadPlanEntry> data_pages;
if (rg_page_index_reader) {
auto offset_index = rg_page_index_reader->GetOffsetIndex(col_idx);
if (offset_index) {
- data_page_filter = MakePageFilter(row_ranges, offset_index,
row_group_row_count);
+ auto row_group_metadata =
reader->metadata()->RowGroup(row_group_index);
+ auto column_chunk = row_group_metadata->ColumnChunk(col_idx);
+ std::optional<DataPageReadPlan> plan = MakeDataPageReadPlan(
+ row_ranges, offset_index, *column_chunk,
row_group_row_count);
+ if (plan) {
+ first_data_page_offset = plan->first_data_page_offset;
+ data_pages = std::move(plan->data_pages);
+ has_data_page_read_plan = true;
+ }
}
}
- return new PageFilteringColumnIterator(col_idx, reader,
std::vector<int>{row_group_index},
- std::move(data_page_filter));
+ // An empty selection still needs the projected Arrow type, especially
for partial nested
+ // projection, but must not construct a PageReader: without a range
cache Arrow eagerly
+ // reads the whole column chunk in GetColumnPageReader(). An iterator
with no row groups
+ // builds the same reader/type tree and immediately reaches EOF
without file I/O.
+ std::vector<int> row_groups;
+ if (!row_ranges.IsEmpty()) {
+ row_groups.push_back(row_group_index);
+ }
+ return new PageFilteringColumnIterator(col_idx, reader,
std::move(row_groups),
+ has_data_page_read_plan,
first_data_page_offset,
+ std::move(data_pages));
};
// Build reader tree with leaf column filtering
@@ -254,14 +340,20 @@ Result<std::shared_ptr<arrow::ChunkedArray>>
PageFilteredRowGroupReader::ReadFil
for (int col_idx : column_reader->LeafColumnIndices()) {
RowRanges effective_ranges = row_ranges;
- int64_t effective_total = row_group_row_count;
- if (rg_page_index_reader) {
+ int64_t effective_total = row_ranges.IsEmpty() ? 0 :
row_group_row_count;
+ if (!row_ranges.IsEmpty() && rg_page_index_reader) {
auto offset_index = rg_page_index_reader->GetOffsetIndex(col_idx);
if (offset_index) {
- auto [compressed, total] =
- ComputeCompressedRowRanges(row_ranges, offset_index,
row_group_row_count);
- effective_ranges = std::move(compressed);
- effective_total = total;
+ auto row_group_metadata =
+
arrow_file_reader->parquet_reader()->metadata()->RowGroup(row_group_index);
+ auto column_chunk = row_group_metadata->ColumnChunk(col_idx);
+ if (MakeDataPageReadPlan(row_ranges, offset_index,
*column_chunk,
+ row_group_row_count)) {
+ auto [compressed, total] =
+ ComputeCompressedRowRanges(row_ranges, offset_index,
row_group_row_count);
+ effective_ranges = std::move(compressed);
+ effective_total = total;
+ }
}
}
@@ -347,6 +439,10 @@ std::vector<::arrow::io::ReadRange>
PageFilteredRowGroupReader::ComputePageRange
const auto& row_ranges = target_row_group.GetRowRanges();
std::vector<::arrow::io::ReadRange> ranges;
+ if (row_ranges.IsEmpty()) {
+ return ranges;
+ }
+
auto file_metadata = parquet_reader->metadata();
auto rg_metadata = file_metadata->RowGroup(row_group_index);
int64_t row_group_row_count = rg_metadata->num_rows();
@@ -359,20 +455,8 @@ std::vector<::arrow::io::ReadRange>
PageFilteredRowGroupReader::ComputePageRange
for (int32_t col_idx : column_indices) {
auto col_chunk = rg_metadata->ColumnChunk(col_idx);
- int64_t data_page_offset = col_chunk->data_page_offset();
- int64_t data_page_compressed_size = col_chunk->total_compressed_size();
- // Dictionary page: always include if present
- if (col_chunk->has_dictionary_page()) {
- int64_t dict_offset = col_chunk->dictionary_page_offset();
- int64_t dict_size = data_page_offset - dict_offset;
- if (dict_size > 0) {
- // if dictionary exists, the data page size should be reduced
by the dictionary
- data_page_compressed_size -= dict_size;
- ranges.push_back({dict_offset, dict_size});
- }
- }
-
- int64_t chunk_end = data_page_offset + data_page_compressed_size;
+ const int64_t column_chunk_offset = GetColumnChunkOffset(*col_chunk);
+ const int64_t column_chunk_compressed_size =
col_chunk->total_compressed_size();
// Try to get OffsetIndex for page-level ranges
std::shared_ptr<::parquet::OffsetIndex> offset_index;
@@ -382,10 +466,27 @@ std::vector<::arrow::io::ReadRange>
PageFilteredRowGroupReader::ComputePageRange
if (!offset_index) {
// No OffsetIndex: fall back to entire column chunk
- ranges.push_back({data_page_offset, data_page_compressed_size});
+ ranges.push_back({column_chunk_offset,
column_chunk_compressed_size});
+ continue;
+ }
+
+ std::optional<DataPageLayout> layout =
+ GetDataPageLayout(*col_chunk, offset_index, row_group_row_count);
+ if (!layout) {
+ // Invalid or empty OffsetIndex: keep the original sequential
reader path.
+ ranges.push_back({column_chunk_offset,
column_chunk_compressed_size});
continue;
}
+ // The full OffsetIndex, rather than data_page_offset in the column
metadata, is the
+ // authoritative location of the first data page. Older parquet-mr
files may omit the
+ // dictionary offset or set it equal to data_page_offset even though a
dictionary prefix
+ // is present (PARQUET-1850/PARQUET-1977).
+ if (layout->first_data_page_offset > layout->column_chunk_offset) {
+ ranges.push_back({layout->column_chunk_offset,
+ layout->first_data_page_offset -
layout->column_chunk_offset});
+ }
+
const auto& page_locations = offset_index->page_locations();
auto num_pages = static_cast<int32_t>(page_locations.size());
@@ -397,11 +498,8 @@ std::vector<::arrow::io::ReadRange>
PageFilteredRowGroupReader::ComputePageRange
continue;
}
- int64_t page_offset = page_locations[page_idx].offset;
- int64_t page_size = (page_idx + 1 < num_pages)
- ? page_locations[page_idx + 1].offset -
page_offset
- : chunk_end - page_offset;
- ranges.push_back({page_offset, page_size});
+ const auto& page = page_locations[page_idx];
+ ranges.push_back({page.offset, page.compressed_page_size});
}
}
diff --git a/src/paimon/format/parquet/page_filtered_row_group_reader.h
b/src/paimon/format/parquet/page_filtered_row_group_reader.h
index 1122ece..98303e7 100644
--- a/src/paimon/format/parquet/page_filtered_row_group_reader.h
+++ b/src/paimon/format/parquet/page_filtered_row_group_reader.h
@@ -20,9 +20,10 @@
#pragma once
#include <cstdint>
-#include <functional>
#include <limits>
#include <memory>
+#include <optional>
+#include <utility>
#include <vector>
#include "arrow/io/caching.h"
@@ -73,6 +74,11 @@ class PageFilteredRowGroupReader {
::parquet::ParquetFileReader* parquet_reader);
private:
+ struct DataPageReadPlan {
+ int64_t first_data_page_offset;
+ std::vector<::parquet::DataPageReadPlanEntry> data_pages;
+ };
+
/// Get the [first_row, last_row] range of a page given page locations.
static std::pair<int64_t, int64_t> GetPageRowRange(
const std::vector<::parquet::PageLocation>& page_locations, int32_t
page_idx,
@@ -87,12 +93,14 @@ class PageFilteredRowGroupReader {
std::shared_ptr<::arrow::MemoryPool> pool,
::parquet::ParquetFileReader*
parquet_reader);
- /// Create a data_page_filter callback for a column based on RowRanges +
OffsetIndex.
- static std::function<bool(const ::parquet::DataPageStats&)> MakePageFilter(
+ /// Build a direct data page read plan for a column based on RowRanges +
OffsetIndex.
+ /// The returned first data page offset and all page offsets are relative
to the
+ /// beginning of the column chunk stream used by Arrow's PageReader.
+ static std::optional<DataPageReadPlan> MakeDataPageReadPlan(
const RowRanges& row_ranges, const
std::shared_ptr<::parquet::OffsetIndex>& offset_index,
- int64_t row_group_row_count);
+ const ::parquet::ColumnChunkMetaData& column_chunk, int64_t
row_group_row_count);
- /// Compute compressed RowRanges after data_page_filter skips non-matching
pages.
+ /// Compute compressed RowRanges after the direct read plan skips
non-matching pages.
static std::pair<RowRanges, int64_t> ComputeCompressedRowRanges(
const RowRanges& original_ranges,
const std::shared_ptr<::parquet::OffsetIndex>& offset_index, int64_t
row_group_row_count);
@@ -103,7 +111,7 @@ class PageFilteredRowGroupReader {
::parquet::arrow::ColumnReader*
column_reader);
/// Read a field (flat or nested) using ColumnReader tree.
- /// Sets data_page_filter on all leaves via factory, then drives each leaf
+ /// Sets a direct page read plan on all leaves via factory, then drives
each leaf
/// independently via ResetLeaf/SkipRecords/ReadRecords using its own
/// compressed_ranges.
static Result<std::shared_ptr<arrow::ChunkedArray>> ReadFilteredField(
diff --git a/src/paimon/format/parquet/page_filtered_row_group_reader_test.cpp
b/src/paimon/format/parquet/page_filtered_row_group_reader_test.cpp
index 2ccdc4e..a00f895 100644
--- a/src/paimon/format/parquet/page_filtered_row_group_reader_test.cpp
+++ b/src/paimon/format/parquet/page_filtered_row_group_reader_test.cpp
@@ -19,10 +19,15 @@
#include "paimon/format/parquet/page_filtered_row_group_reader.h"
+#include <algorithm>
#include <cstdint>
+#include <cstring>
+#include <functional>
#include <iostream>
+#include <limits>
#include <map>
#include <memory>
+#include <mutex>
#include <optional>
#include <string>
#include <utility>
@@ -32,10 +37,12 @@
#include "arrow/array/array_nested.h"
#include "arrow/c/abi.h"
#include "arrow/c/bridge.h"
+#include "arrow/io/api.h"
#include "arrow/ipc/json_simple.h"
#include "gtest/gtest.h"
#include "paimon/common/utils/arrow/arrow_input_stream_adapter.h"
#include "paimon/common/utils/arrow/mem_utils.h"
+#include "paimon/common/utils/stream_utils.h"
#include "paimon/defs.h"
#include "paimon/format/parquet/parquet_file_batch_reader.h"
#include "paimon/format/parquet/parquet_format_defs.h"
@@ -50,6 +57,7 @@
#include "paimon/testing/utils/testharness.h"
#include "paimon/utils/roaring_bitmap32.h"
#include "parquet/arrow/reader.h"
+#include "parquet/column_page.h"
#include "parquet/file_reader.h"
#include "parquet/properties.h"
@@ -59,6 +67,67 @@ class Predicate;
namespace paimon::parquet::test {
+class ReadAtTrackingInputStream : public InputStream {
+ public:
+ explicit ReadAtTrackingInputStream(std::shared_ptr<InputStream> input)
+ : input_(std::move(input)) {}
+
+ Status Seek(int64_t offset, SeekOrigin origin) override {
+ return input_->Seek(offset, origin);
+ }
+
+ Result<int64_t> GetPos() const override {
+ return input_->GetPos();
+ }
+
+ Result<int64_t> Read(char* buffer, int64_t size) override {
+ return input_->Read(buffer, size);
+ }
+
+ Result<int64_t> Read(char* buffer, int64_t size, int64_t offset) override {
+ RecordPositionalRead(offset, size);
+ return input_->Read(buffer, size, offset);
+ }
+
+ void ReadAsync(char* buffer, int64_t size, int64_t offset,
+ std::function<void(Status)>&& callback) override {
+ RecordPositionalRead(offset, size);
+ input_->ReadAsync(buffer, size, offset, std::move(callback));
+ }
+
+ Result<std::string> GetUri() const override {
+ return input_->GetUri();
+ }
+
+ Result<int64_t> Length() const override {
+ return input_->Length();
+ }
+
+ Status Close() override {
+ return input_->Close();
+ }
+
+ void ClearReadAtRanges() {
+ std::lock_guard<std::mutex> lock(mutex_);
+ read_at_ranges_.clear();
+ }
+
+ std::vector<arrow::io::ReadRange> GetReadAtRanges() const {
+ std::lock_guard<std::mutex> lock(mutex_);
+ return read_at_ranges_;
+ }
+
+ private:
+ void RecordPositionalRead(int64_t offset, int64_t size) {
+ std::lock_guard<std::mutex> lock(mutex_);
+ read_at_ranges_.push_back({offset, size});
+ }
+
+ std::shared_ptr<InputStream> input_;
+ mutable std::mutex mutex_;
+ std::vector<arrow::io::ReadRange> read_at_ranges_;
+};
+
/// Test fixture for page-level filtering.
/// Creates Parquet files with multiple row groups and small page sizes to
ensure
/// multiple pages per row group, enabling page-level filtering tests.
@@ -77,10 +146,12 @@ class PageFilteredRowGroupReaderTest : public
::testing::Test {
/// @param struct_array Data to write
/// @param write_batch_size Controls page size (number of rows per page)
/// @param max_row_group_length Controls row group size
- void WriteTestFile(const std::string& file_name,
- const std::shared_ptr<arrow::StructArray>& struct_array,
- int32_t write_batch_size, int64_t max_row_group_length,
- bool enable_dictionary = false, int64_t data_page_size
= 1) {
+ void WriteTestFile(
+ const std::string& file_name, const
std::shared_ptr<arrow::StructArray>& struct_array,
+ int32_t write_batch_size, int64_t max_row_group_length, bool
enable_dictionary = false,
+ int64_t data_page_size = 1, bool enable_page_index = true,
+ ::parquet::ParquetDataPageVersion data_page_version =
::parquet::ParquetDataPageVersion::V1,
+ int64_t dictionary_page_size_limit = -1) {
auto data_type = struct_array->struct_type();
auto data_schema = arrow::schema(data_type->fields());
auto data_arrow_array = std::make_unique<ArrowArray>();
@@ -92,10 +163,18 @@ class PageFilteredRowGroupReaderTest : public
::testing::Test {
builder.max_row_group_length(max_row_group_length);
if (enable_dictionary) {
builder.enable_dictionary();
+ if (dictionary_page_size_limit >= 0) {
+ builder.dictionary_pagesize_limit(dictionary_page_size_limit);
+ }
} else {
builder.disable_dictionary(); // Ensure page index min/max are
meaningful
}
- builder.enable_write_page_index(); // Enable page index for
page-level filtering
+ if (enable_page_index) {
+ builder.enable_write_page_index();
+ } else {
+ builder.disable_write_page_index();
+ }
+ builder.data_page_version(data_page_version);
// Data page size controls when a page is flushed. The default of 1
byte forces a new
// page after every write_batch_size rows (each batch becomes one
page), giving pages
// aligned across columns. A larger byte-based value combined with
write_batch_size=1
@@ -138,15 +217,19 @@ class PageFilteredRowGroupReaderTest : public
::testing::Test {
}
/// Read back a Parquet file with a predicate, a bitmap, and page index
filter enabled.
- void ReadWithPredicateAndBitmapImpl(const std::string& file_name,
- const std::shared_ptr<arrow::Schema>&
read_schema,
- const std::shared_ptr<Predicate>&
predicate,
- const RoaringBitmap32& bitmap,
- std::shared_ptr<arrow::ChunkedArray>*
out,
- const std::map<std::string,
std::string> options = {},
- int32_t batch_size = 1024) {
+ void ReadWithPredicateAndBitmapImpl(
+ const std::string& file_name, const std::shared_ptr<arrow::Schema>&
read_schema,
+ const std::shared_ptr<Predicate>& predicate, const RoaringBitmap32&
bitmap,
+ std::shared_ptr<arrow::ChunkedArray>* out,
+ const std::map<std::string, std::string> options = {}, int32_t
batch_size = 1024,
+ std::vector<arrow::io::ReadRange>* read_at_ranges = nullptr) {
ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ std::shared_ptr<ReadAtTrackingInputStream> tracking_in;
+ if (read_at_ranges != nullptr) {
+ tracking_in = std::make_shared<ReadAtTrackingInputStream>(in);
+ in = tracking_in;
+ }
auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
ASSERT_OK_AND_ASSIGN(
@@ -156,8 +239,16 @@ class PageFilteredRowGroupReaderTest : public
::testing::Test {
auto c_schema = std::make_unique<ArrowSchema>();
ASSERT_TRUE(arrow::ExportSchema(*read_schema, c_schema.get()).ok());
ASSERT_OK(batch_reader->SetReadSchema(c_schema.get(), predicate,
bitmap));
+ if (tracking_in) {
+ // Ignore footer and schema metadata reads. The test only checks
I/O issued while
+ // consuming data pages.
+ tracking_in->ClearReadAtRanges();
+ }
ASSERT_OK_AND_ASSIGN(*out,
paimon::test::ReadResultCollector::CollectResult(batch_reader.get()));
+ if (tracking_in) {
+ *read_at_ranges = tracking_in->GetReadAtRanges();
+ }
}
protected:
@@ -197,6 +288,50 @@ static std::shared_ptr<arrow::StructArray>
MakeTwoColumnData(int32_t num_rows) {
return arrow::StructArray::Make({a_array, b_array}, {field_a,
field_b}).ValueOrDie();
}
+static void AssertUnselectedPageHeadersNotRead(
+ const std::vector<::parquet::PageLocation>& page_locations,
+ const std::vector<int32_t>& selected_pages,
+ const std::vector<arrow::io::ReadRange>& read_at_ranges) {
+ for (int32_t page_idx = 0; page_idx <
static_cast<int32_t>(page_locations.size()); ++page_idx) {
+ if (std::find(selected_pages.begin(), selected_pages.end(), page_idx)
!=
+ selected_pages.end()) {
+ continue;
+ }
+ const int64_t page_header_offset = page_locations[page_idx].offset;
+ for (const auto& read_range : read_at_ranges) {
+ ASSERT_FALSE(read_range.offset <= page_header_offset &&
+ page_header_offset < read_range.offset +
read_range.length)
+ << "unselected page " << page_idx << " header at " <<
page_header_offset
+ << " was covered by positional ReadAt [" << read_range.offset
<< ", "
+ << read_range.offset + read_range.length << ")";
+ }
+ }
+}
+
+static void AssertReadRangeCovered(const std::vector<arrow::io::ReadRange>&
read_ranges,
+ int64_t expected_offset, int64_t
expected_size) {
+ ASSERT_GT(expected_size, 0);
+ ASSERT_TRUE(std::any_of(
+ read_ranges.begin(), read_ranges.end(),
+ [expected_offset, expected_size](const arrow::io::ReadRange&
read_range) {
+ return read_range.offset <= expected_offset &&
+ expected_offset + expected_size <= read_range.offset +
read_range.length;
+ }))
+ << "expected range [" << expected_offset << ", " << expected_offset +
expected_size
+ << ") was not covered by any positional read";
+}
+
+static void AssertSelectedPagesRead(const
std::vector<::parquet::PageLocation>& page_locations,
+ const std::vector<int32_t>& selected_pages,
+ const std::vector<arrow::io::ReadRange>&
read_ranges) {
+ for (int32_t page_idx : selected_pages) {
+ ASSERT_GE(page_idx, 0);
+ ASSERT_LT(page_idx, static_cast<int32_t>(page_locations.size()));
+ const auto& page = page_locations[page_idx];
+ AssertReadRangeCovered(read_ranges, page.offset,
page.compressed_page_size);
+ }
+}
+
/// Test: page-level filtering correctly skips non-matching pages.
///
/// Scenario: 100 rows, 10 rows per page, 1 row group.
@@ -1183,6 +1318,424 @@ TEST_F(PageFilteredRowGroupReaderTest,
BitmapAllPagesSomeRowGroups) {
}
}
+/// Test: OffsetIndex lets the reader jump directly to selected data pages.
+///
+/// The bitmap selects rows from page 1 and page 8. Synchronous positional
reads issued while
+/// consuming the row group must not cover any unselected page header. Before
direct page jumps,
+/// SerializedPageReader sequentially Peeked every page header and this
assertion failed.
+TEST_F(PageFilteredRowGroupReaderTest,
DirectOffsetIndexJumpDoesNotReadUnselectedPageHeaders) {
+ std::string file_name = dir_->Str() + "/direct_offset_index_jump.parquet";
+ auto data = MakeSequentialIntData(100);
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/100);
+
+ std::vector<::parquet::PageLocation> page_locations;
+ {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
+ ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
+ auto parquet_reader = ::parquet::ParquetFileReader::Open(in_stream);
+ ASSERT_TRUE(parquet_reader);
+
+ auto page_index_reader = parquet_reader->GetPageIndexReader();
+ ASSERT_TRUE(page_index_reader);
+ auto row_group_page_index = page_index_reader->RowGroup(0);
+ ASSERT_TRUE(row_group_page_index);
+ auto offset_index = row_group_page_index->GetOffsetIndex(0);
+ ASSERT_TRUE(offset_index);
+ page_locations = offset_index->page_locations();
+ }
+
+ ASSERT_EQ(10, page_locations.size());
+ for (int32_t page_idx = 0; page_idx < 10; ++page_idx) {
+ ASSERT_EQ(page_idx * 10, page_locations[page_idx].first_row_index);
+ }
+
+ RoaringBitmap32 bitmap;
+ bitmap.Add(15); // page 1
+ bitmap.Add(85); // page 8
+
+ std::map<std::string, std::string> options;
+ options[PARQUET_READ_ROW_RANGES_COALESCE_HOLE_SIZE_LIMIT] = "0";
+ options[PARQUET_READ_CACHE_OPTION_HOLE_SIZE_LIMIT] = "0";
+
+ auto read_schema = arrow::schema({arrow::field("val", arrow::int32())});
+ std::shared_ptr<arrow::ChunkedArray> result;
+ std::vector<arrow::io::ReadRange> read_at_ranges;
+ ReadWithPredicateAndBitmapImpl(file_name, read_schema,
/*predicate=*/nullptr, bitmap, &result,
+ options, /*batch_size=*/1024,
&read_at_ranges);
+
+ ASSERT_TRUE(result);
+ ASSERT_EQ(2, result->length());
+ auto flat = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto struct_arr = std::dynamic_pointer_cast<arrow::StructArray>(flat);
+ ASSERT_TRUE(struct_arr);
+ auto val_arr =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(0));
+ ASSERT_TRUE(val_arr);
+ ASSERT_EQ(15, val_arr->Value(0));
+ ASSERT_EQ(85, val_arr->Value(1));
+
+ AssertSelectedPagesRead(page_locations, /*selected_pages=*/{1, 8},
read_at_ranges);
+ AssertUnselectedPageHeadersNotRead(page_locations, /*selected_pages=*/{1,
8}, read_at_ranges);
+}
+
+/// Dictionary pages are column-local. Each leaf must load its own dictionary
before the
+/// PageReader jumps directly to a selected late data page.
+TEST_F(PageFilteredRowGroupReaderTest,
DirectOffsetIndexJumpReadsEachLeafDictionary) {
+ std::string file_name = dir_->Str() +
"/direct_offset_index_jump_dictionary.parquet";
+
+ arrow::Int32Builder a_builder;
+ arrow::Int32Builder b_builder;
+ ASSERT_TRUE(a_builder.Reserve(100).ok());
+ ASSERT_TRUE(b_builder.Reserve(100).ok());
+ for (int32_t i = 0; i < 100; ++i) {
+ a_builder.UnsafeAppend(i % 7);
+ b_builder.UnsafeAppend(i % 5);
+ }
+ auto a_array = a_builder.Finish().ValueOrDie();
+ auto b_array = b_builder.Finish().ValueOrDie();
+ auto field_a = arrow::field("a", arrow::int32());
+ auto field_b = arrow::field("b", arrow::int32());
+ auto data = arrow::StructArray::Make({a_array, b_array}, {field_a,
field_b}).ValueOrDie();
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/100,
+ /*enable_dictionary=*/true);
+
+ std::vector<std::vector<::parquet::PageLocation>> page_locations(2);
+ std::vector<arrow::io::ReadRange> dictionary_ranges;
+ {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
+ ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
+ auto parquet_reader = ::parquet::ParquetFileReader::Open(in_stream);
+ ASSERT_TRUE(parquet_reader);
+
+ auto row_group = parquet_reader->metadata()->RowGroup(0);
+ auto page_index_reader = parquet_reader->GetPageIndexReader();
+ ASSERT_TRUE(page_index_reader);
+ auto row_group_page_index = page_index_reader->RowGroup(0);
+ ASSERT_TRUE(row_group_page_index);
+
+ RowRanges selected_rows;
+ selected_rows.Add(RowRanges::Range(95, 95));
+ auto ranges = PageFilteredRowGroupReader::ComputePageRanges(
+ TargetRowGroup(/*rg_index=*/0, /*is_partially_matched=*/true,
+ /*ranges=*/selected_rows),
+ /*column_indices=*/{0, 1}, parquet_reader.get());
+
+ for (int32_t col_idx = 0; col_idx < 2; ++col_idx) {
+ auto column_chunk = row_group->ColumnChunk(col_idx);
+ ASSERT_TRUE(column_chunk->has_dictionary_page());
+ const int64_t dictionary_offset =
column_chunk->dictionary_page_offset();
+ const int64_t data_page_offset = column_chunk->data_page_offset();
+ ASSERT_LT(dictionary_offset, data_page_offset);
+
+ auto offset_index = row_group_page_index->GetOffsetIndex(col_idx);
+ ASSERT_TRUE(offset_index);
+ page_locations[col_idx] = offset_index->page_locations();
+ ASSERT_EQ(10, page_locations[col_idx].size());
+
+ const auto& selected_page = page_locations[col_idx][9];
+ const auto contains_range = [&ranges](int64_t offset, int64_t
size) {
+ return std::any_of(ranges.begin(), ranges.end(),
+ [offset, size](const arrow::io::ReadRange&
range) {
+ return range.offset == offset &&
range.length == size;
+ });
+ };
+ ASSERT_TRUE(contains_range(dictionary_offset, data_page_offset -
dictionary_offset));
+ ASSERT_TRUE(contains_range(selected_page.offset,
selected_page.compressed_page_size));
+ dictionary_ranges.push_back({dictionary_offset, data_page_offset -
dictionary_offset});
+
+ auto page_reader =
parquet_reader->RowGroup(0)->GetColumnPageReader(col_idx);
+ ASSERT_TRUE(page_reader);
+ std::vector<::parquet::Encoding::type> data_page_encodings;
+ while (std::shared_ptr<::parquet::Page> page =
page_reader->NextPage()) {
+ if (page->type() == ::parquet::PageType::DATA_PAGE ||
+ page->type() == ::parquet::PageType::DATA_PAGE_V2) {
+ data_page_encodings.push_back(
+
std::static_pointer_cast<::parquet::DataPage>(page)->encoding());
+ }
+ }
+ ASSERT_EQ(page_locations[col_idx].size(),
data_page_encodings.size());
+ ASSERT_TRUE(data_page_encodings[9] ==
::parquet::Encoding::PLAIN_DICTIONARY ||
+ data_page_encodings[9] ==
::parquet::Encoding::RLE_DICTIONARY);
+ }
+ }
+
+ RoaringBitmap32 bitmap;
+ bitmap.Add(95); // Last data page; every leaf still needs its dictionary
first.
+
+ std::map<std::string, std::string> options;
+ options[PARQUET_READ_ROW_RANGES_COALESCE_HOLE_SIZE_LIMIT] = "0";
+ options[PARQUET_READ_CACHE_OPTION_HOLE_SIZE_LIMIT] = "0";
+
+ auto read_schema = arrow::schema({field_a, field_b});
+ std::shared_ptr<arrow::ChunkedArray> result;
+ std::vector<arrow::io::ReadRange> read_at_ranges;
+ ReadWithPredicateAndBitmapImpl(file_name, read_schema,
/*predicate=*/nullptr, bitmap, &result,
+ options, /*batch_size=*/1024,
&read_at_ranges);
+
+ ASSERT_TRUE(result);
+ ASSERT_EQ(1, result->length());
+ auto flat = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto struct_arr = std::dynamic_pointer_cast<arrow::StructArray>(flat);
+ ASSERT_TRUE(struct_arr);
+ auto result_a =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(0));
+ auto result_b =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(1));
+ ASSERT_TRUE(result_a);
+ ASSERT_TRUE(result_b);
+ ASSERT_EQ(95 % 7, result_a->Value(0));
+ ASSERT_EQ(95 % 5, result_b->Value(0));
+
+ for (size_t col_idx = 0; col_idx < page_locations.size(); ++col_idx) {
+ AssertReadRangeCovered(read_at_ranges,
dictionary_ranges[col_idx].offset,
+ dictionary_ranges[col_idx].length);
+ AssertSelectedPagesRead(page_locations[col_idx],
/*selected_pages=*/{9}, read_at_ranges);
+ const auto& leaf_page_locations = page_locations[col_idx];
+ AssertUnselectedPageHeadersNotRead(leaf_page_locations,
/*selected_pages=*/{9},
+ read_at_ranges);
+ }
+}
+
+/// A column can switch from dictionary encoding to PLAIN after the dictionary
reaches its size
+/// limit. A single direct read plan must decode both kinds of selected data
pages correctly.
+TEST_F(PageFilteredRowGroupReaderTest,
DirectOffsetIndexJumpSupportsDictionaryFallbackToPlain) {
+ std::string file_name = dir_->Str() +
"/direct_offset_index_dictionary_fallback.parquet";
+ auto data = MakeSequentialIntData(1000);
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/1000,
+ /*enable_dictionary=*/true, /*data_page_size=*/1,
+ /*enable_page_index=*/true,
::parquet::ParquetDataPageVersion::V1,
+ /*dictionary_page_size_limit=*/256);
+
+ std::vector<::parquet::PageLocation> page_locations;
+ std::vector<::parquet::Encoding::type> data_page_encodings;
+ int64_t dictionary_offset = 0;
+ int64_t first_data_page_offset = 0;
+ {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
+ ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
+ auto parquet_reader = ::parquet::ParquetFileReader::Open(in_stream);
+ ASSERT_TRUE(parquet_reader);
+
+ auto column_chunk =
parquet_reader->metadata()->RowGroup(0)->ColumnChunk(0);
+ ASSERT_TRUE(column_chunk->has_dictionary_page());
+ dictionary_offset = column_chunk->dictionary_page_offset();
+
+ auto page_index_reader = parquet_reader->GetPageIndexReader();
+ ASSERT_TRUE(page_index_reader);
+ auto row_group_page_index = page_index_reader->RowGroup(0);
+ ASSERT_TRUE(row_group_page_index);
+ auto offset_index = row_group_page_index->GetOffsetIndex(0);
+ ASSERT_TRUE(offset_index);
+ page_locations = offset_index->page_locations();
+ ASSERT_FALSE(page_locations.empty());
+ first_data_page_offset = page_locations.front().offset;
+ ASSERT_LT(dictionary_offset, first_data_page_offset);
+
+ auto page_reader = parquet_reader->RowGroup(0)->GetColumnPageReader(0);
+ ASSERT_TRUE(page_reader);
+ while (std::shared_ptr<::parquet::Page> page =
page_reader->NextPage()) {
+ if (page->type() == ::parquet::PageType::DATA_PAGE ||
+ page->type() == ::parquet::PageType::DATA_PAGE_V2) {
+ data_page_encodings.push_back(
+
std::static_pointer_cast<::parquet::DataPage>(page)->encoding());
+ }
+ }
+ }
+ ASSERT_EQ(page_locations.size(), data_page_encodings.size());
+
+ int32_t dictionary_page_idx = -1;
+ int32_t plain_page_idx = -1;
+ for (int32_t page_idx = 0; page_idx <
static_cast<int32_t>(data_page_encodings.size());
+ ++page_idx) {
+ const auto encoding = data_page_encodings[page_idx];
+ if (dictionary_page_idx < 0 && (encoding ==
::parquet::Encoding::PLAIN_DICTIONARY ||
+ encoding ==
::parquet::Encoding::RLE_DICTIONARY)) {
+ dictionary_page_idx = page_idx;
+ }
+ if (encoding == ::parquet::Encoding::PLAIN) {
+ plain_page_idx = page_idx;
+ }
+ }
+ ASSERT_GE(dictionary_page_idx, 0);
+ ASSERT_GT(plain_page_idx, dictionary_page_idx);
+
+ const auto dictionary_row =
+
static_cast<int32_t>(page_locations[dictionary_page_idx].first_row_index);
+ const auto plain_row =
static_cast<int32_t>(page_locations[plain_page_idx].first_row_index);
+ RoaringBitmap32 bitmap;
+ bitmap.Add(dictionary_row);
+ bitmap.Add(plain_row);
+
+ std::map<std::string, std::string> options;
+ options[PARQUET_READ_ROW_RANGES_COALESCE_HOLE_SIZE_LIMIT] = "0";
+ options[PARQUET_READ_CACHE_OPTION_HOLE_SIZE_LIMIT] = "0";
+
+ auto read_schema = arrow::schema({arrow::field("val", arrow::int32())});
+ std::shared_ptr<arrow::ChunkedArray> result;
+ std::vector<arrow::io::ReadRange> read_at_ranges;
+ ReadWithPredicateAndBitmapImpl(file_name, read_schema,
/*predicate=*/nullptr, bitmap, &result,
+ options, /*batch_size=*/1024,
&read_at_ranges);
+
+ ASSERT_TRUE(result);
+ ASSERT_EQ(2, result->length());
+ auto flat = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto struct_arr = std::dynamic_pointer_cast<arrow::StructArray>(flat);
+ ASSERT_TRUE(struct_arr);
+ auto values =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(0));
+ ASSERT_TRUE(values);
+ ASSERT_EQ(dictionary_row, values->Value(0));
+ ASSERT_EQ(plain_row, values->Value(1));
+
+ AssertReadRangeCovered(read_at_ranges, dictionary_offset,
+ first_data_page_offset - dictionary_offset);
+ AssertSelectedPagesRead(page_locations, {dictionary_page_idx,
plain_page_idx}, read_at_ranges);
+ AssertUnselectedPageHeadersNotRead(page_locations, {dictionary_page_idx,
plain_page_idx},
+ read_at_ranges);
+}
+
+/// Empty row selection may read OffsetIndex metadata, but must not read
dictionary/data pages.
+TEST_F(PageFilteredRowGroupReaderTest,
DictionaryEmptySelectionDoesNotReadPages) {
+ std::string file_name = dir_->Str() + "/dictionary_empty_bitmap.parquet";
+ auto data = MakeSequentialIntData(100);
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/100,
+ /*enable_dictionary=*/true);
+
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> input,
fs_->Open(file_name));
+ auto tracking_input = std::make_shared<ReadAtTrackingInputStream>(input);
+ ASSERT_OK_AND_ASSIGN(int64_t length, tracking_input->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(tracking_input,
length, arrow_pool_);
+
+ ::parquet::arrow::FileReaderBuilder builder;
+ ASSERT_TRUE(builder.Open(in_stream).ok());
+ builder.memory_pool(arrow_pool_.get());
+ auto arrow_file_reader_result = builder.Build();
+ ASSERT_TRUE(arrow_file_reader_result.ok()) <<
arrow_file_reader_result.status().ToString();
+ std::unique_ptr<::parquet::arrow::FileReader> arrow_file_reader =
+ std::move(arrow_file_reader_result).ValueOrDie();
+
+ auto column_chunk =
+
arrow_file_reader->parquet_reader()->metadata()->RowGroup(0)->ColumnChunk(0);
+ ASSERT_TRUE(column_chunk->has_dictionary_page());
+ const int64_t column_chunk_offset = column_chunk->dictionary_page_offset();
+ const int64_t column_chunk_end = column_chunk_offset +
column_chunk->total_compressed_size();
+ ASSERT_GE(column_chunk_offset, 0);
+ ASSERT_GT(column_chunk_end, column_chunk_offset);
+
+ RowRanges empty_ranges;
+ TargetRowGroup empty_target(/*rg_index=*/0, /*is_partially_matched=*/true,
+ /*ranges=*/empty_ranges);
+ auto page_ranges = PageFilteredRowGroupReader::ComputePageRanges(
+ empty_target, /*column_indices=*/{0},
arrow_file_reader->parquet_reader());
+ ASSERT_TRUE(page_ranges.empty());
+
+ tracking_input->ClearReadAtRanges();
+ ASSERT_OK_AND_ASSIGN(
+ std::unique_ptr<arrow::RecordBatchReader> result_reader,
+ PageFilteredRowGroupReader::ReadFilteredRowGroup(
+ empty_target, /*column_indices=*/{0},
arrow::io::CacheOptions::Defaults(),
+ /*pre_buffered=*/false, /*page_ranges=*/{},
/*max_chunksize=*/1024, arrow_pool_,
+ arrow_file_reader.get()));
+ std::shared_ptr<arrow::RecordBatch> batch;
+ ASSERT_TRUE(result_reader->ReadNext(&batch).ok());
+ ASSERT_FALSE(batch);
+ for (const auto& read_range : tracking_input->GetReadAtRanges()) {
+ const int64_t read_end = read_range.offset + read_range.length;
+ ASSERT_TRUE(read_end <= column_chunk_offset || read_range.offset >=
column_chunk_end)
+ << "empty selection read column chunk [" << column_chunk_offset <<
", "
+ << column_chunk_end << ") via positional range [" <<
read_range.offset << ", "
+ << read_end << ")";
+ }
+}
+
+/// The direct page plan must work for DATA_PAGE_V2 as well as DATA_PAGE_V1.
Selecting adjacent
+/// pages also verifies that the plan cursor advances exactly once per page.
+TEST_F(PageFilteredRowGroupReaderTest,
DirectOffsetIndexJumpDataPageV2AdjacentPages) {
+ std::string file_name = dir_->Str() +
"/direct_offset_index_jump_v2.parquet";
+ auto data = MakeSequentialIntData(100);
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/100,
+ /*enable_dictionary=*/false, /*data_page_size=*/1,
+ /*enable_page_index=*/true,
::parquet::ParquetDataPageVersion::V2);
+
+ std::vector<::parquet::PageLocation> page_locations;
+ {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
+ ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
+ auto parquet_reader = ::parquet::ParquetFileReader::Open(in_stream);
+ auto page_index_reader = parquet_reader->GetPageIndexReader();
+ ASSERT_TRUE(page_index_reader);
+ auto row_group_page_index = page_index_reader->RowGroup(0);
+ ASSERT_TRUE(row_group_page_index);
+ auto offset_index = row_group_page_index->GetOffsetIndex(0);
+ ASSERT_TRUE(offset_index);
+ page_locations = offset_index->page_locations();
+ }
+ ASSERT_EQ(10, page_locations.size());
+
+ RoaringBitmap32 bitmap;
+ bitmap.Add(45); // Page 4.
+ bitmap.Add(55); // Adjacent page 5.
+
+ std::map<std::string, std::string> options;
+ options[PARQUET_READ_ROW_RANGES_COALESCE_HOLE_SIZE_LIMIT] = "0";
+ options[PARQUET_READ_CACHE_OPTION_HOLE_SIZE_LIMIT] = "0";
+
+ auto read_schema = arrow::schema({arrow::field("val", arrow::int32())});
+ std::shared_ptr<arrow::ChunkedArray> result;
+ std::vector<arrow::io::ReadRange> read_at_ranges;
+ ReadWithPredicateAndBitmapImpl(file_name, read_schema,
/*predicate=*/nullptr, bitmap, &result,
+ options, /*batch_size=*/1024,
&read_at_ranges);
+
+ ASSERT_TRUE(result);
+ ASSERT_EQ(2, result->length());
+ auto flat = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto struct_arr = std::dynamic_pointer_cast<arrow::StructArray>(flat);
+ ASSERT_TRUE(struct_arr);
+ auto values =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(0));
+ ASSERT_TRUE(values);
+ ASSERT_EQ(45, values->Value(0));
+ ASSERT_EQ(55, values->Value(1));
+ AssertSelectedPagesRead(page_locations, /*selected_pages=*/{4, 5},
read_at_ranges);
+ AssertUnselectedPageHeadersNotRead(page_locations, /*selected_pages=*/{4,
5}, read_at_ranges);
+}
+
+/// Files without OffsetIndex must keep the original full-column reader path.
+TEST_F(PageFilteredRowGroupReaderTest,
MissingOffsetIndexFallsBackToSequentialRead) {
+ std::string file_name = dir_->Str() +
"/missing_offset_index_fallback.parquet";
+ auto data = MakeSequentialIntData(100);
+ WriteTestFile(file_name, data, /*write_batch_size=*/10,
/*max_row_group_length=*/100,
+ /*enable_dictionary=*/false, /*data_page_size=*/1,
+ /*enable_page_index=*/false);
+
+ {
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> in,
fs_->Open(file_name));
+ ASSERT_OK_AND_ASSIGN(int64_t length, in->Length());
+ auto in_stream = std::make_shared<ArrowInputStreamAdapter>(in, length,
arrow_pool_);
+ auto parquet_reader = ::parquet::ParquetFileReader::Open(in_stream);
+ auto page_index_reader = parquet_reader->GetPageIndexReader();
+ ASSERT_TRUE(page_index_reader);
+ ASSERT_FALSE(page_index_reader->RowGroup(0));
+ }
+
+ RoaringBitmap32 bitmap;
+ bitmap.Add(15);
+ bitmap.Add(85);
+
+ auto read_schema = arrow::schema({arrow::field("val", arrow::int32())});
+ std::shared_ptr<arrow::ChunkedArray> result;
+ ReadWithPredicateAndBitmapImpl(file_name, read_schema,
/*predicate=*/nullptr, bitmap, &result);
+
+ ASSERT_TRUE(result);
+ ASSERT_EQ(2, result->length());
+ auto flat = arrow::Concatenate(result->chunks()).ValueOrDie();
+ auto struct_arr = std::dynamic_pointer_cast<arrow::StructArray>(flat);
+ ASSERT_TRUE(struct_arr);
+ auto values =
std::dynamic_pointer_cast<arrow::Int32Array>(struct_arr->field(0));
+ ASSERT_TRUE(values);
+ ASSERT_EQ(15, values->Value(0));
+ ASSERT_EQ(85, values->Value(1));
+}
+
/// Test: bitmap hits partial pages of a row group (no predicate).
///
/// 200 rows, 10 rows per page, 100 rows per row group → 2 row groups.