This is an automated email from the ASF dual-hosted git repository.
Gabriel39 pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new cc2d6a82a57 [improvement](parquet) Support compound Page Index pruning
in File Scanner V2 (#66412)
cc2d6a82a57 is described below
commit cc2d6a82a57d7adfdc8bd00313e96f466a32f9db
Author: Gabriel <[email protected]>
AuthorDate: Wed Aug 5 18:01:18 2026 +0800
[improvement](parquet) Support compound Page Index pruning in File Scanner
V2 (#66412)
## What problem does this PR solve?
File Scanner V2 can use full multi-column predicates for Row Group
statistics, but Page Index
pruning previously handled only predicates that referenced one slot.
Compound predicates such as
the multi-column `OR` expressions in TPC-DS Q28 and the DNF expressions
in Q88 therefore decoded
the complete Row Group before residual filtering.
## What is changed?
- Evaluate multi-column `AND`/`OR` predicate trees against Parquet Page
Index metadata.
- Intersect child candidate row ranges for `AND` and union them for
`OR`, following the predicate
tree rather than constructing decoder-level branch masks.
- Cache decoded page zone maps per slot while evaluating multiple leaves
from the same column.
- Fall back conservatively: an unavailable `AND` child contributes no
pruning, while an unavailable
`OR` branch retains the complete Row Group range.
- Keep the complete compound expression as a row-level residual
predicate after all referenced
columns are materialized, preserving SQL NULL and expression semantics.
- Preserve the existing single-column dictionary, raw-filter,
lazy-materialization, and Page Index
paths unchanged.
- Remove the auxiliary-reader and decoder branch-mask implementation
from the earlier revision.
## Performance
The Release microbenchmark scans the same 1,048,576-row PLAIN fixture in
both variants and changes
only the Doris Page Index switch. Both variants evaluate the same
two-column `OR` residual and
return exactly 209,712 rows. Results were pinned to one CPU, preceded by
three warmups, and collected
in ABBA order with 10 repetitions and a 0.5-second minimum per
repetition.
| Variant | Median CPU time | CPU time per raw row | Selected rows |
|---|---:|---:|---:|
| Page Index off | 12.51 ms | 11.93 ns | 209,712 |
| Page Index on | 9.05 ms | 8.63 ns | 209,712 |
| **Change** | **-27.67%** | **-27.67%** | **identical** |
The single-column Page Index path is structurally unchanged; only
conjuncts that do not resolve to
a single slot enter the new evaluator.
## Test
- 17 focused ASAN tests passed for native Parquet statistics and scanner
behavior.
- Coverage includes multi-column OR union, AND intersection, nested DNF,
missing-index fallback,
exact residual results, and existing single-column Page Index/range-gap
pruning.
- 13 Parquet benchmark scenario tests passed.
- Release benchmark build succeeded and all 169 `ParquetReader` cases
were registered.
- The dedicated Page Index on/off benchmark completed without errors or
row-count mismatches.
---
be/benchmark/parquet/AGENTS.md | 8 +-
be/benchmark/parquet/README.md | 12 ++
be/benchmark/parquet/benchmark_parquet_reader.hpp | 169 +++++++++++++++++
be/src/format_v2/parquet/parquet_statistics.cpp | 208 ++++++++++++++++++++-
be/test/format_v2/parquet/parquet_scan_test.cpp | 57 ++++++
.../format_v2/parquet/parquet_statistics_test.cpp | 160 ++++++++++++++++
docs/file-scanner-v2-parquet-scan-design.md | 17 +-
7 files changed, 616 insertions(+), 15 deletions(-)
diff --git a/be/benchmark/parquet/AGENTS.md b/be/benchmark/parquet/AGENTS.md
index 4d8c0610f17..da513488520 100644
--- a/be/benchmark/parquet/AGENTS.md
+++ b/be/benchmark/parquet/AGENTS.md
@@ -57,7 +57,7 @@ be/output/lib/benchmark_test --benchmark_list_tests \
| grep -c '^ParquetSelection/' # currently 25
be/output/lib/benchmark_test --benchmark_list_tests \
- | grep -c '^ParquetReader/' # currently 167
+ | grep -c '^ParquetReader/' # currently 169
be/output/lib/benchmark_test --benchmark_list_tests \
| grep -c '^FileScannerExpr/' # currently 8
@@ -189,6 +189,10 @@ Except for the axis being varied, reader cases inherit the
baseline: nullable IN
alternating 10% nulls, 10% selectivity, 32 columns, predicate at column zero,
and predicate plus
payload projection.
+Two dedicated multi-column OR cases scan the same Page Index fixture and
change only the Doris Page
+Index switch. They retain the complete residual expression and validate the
same selected row
+count before reporting throughput.
+
## How decoder data is generated
Decoder pages are constructed in memory before the timed loop. There is no
Parquet file, Python
@@ -346,7 +350,7 @@ be simulated by silently changing the local reader
benchmark.
## Current validation record
-The current expected registration counts are 228 decoder, 92 kernel, 25
selection, 167 reader, and
+The current expected registration counts are 228 decoder, 92 kernel, 25
selection, 169 reader, and
8 expression-lifecycle cases. A smoke run is an execution record only, not a
reviewed performance
baseline, because repetitions, host isolation, warmups, cache control, `perf`
data, variance, and
before/after comparison are not collected.
diff --git a/be/benchmark/parquet/README.md b/be/benchmark/parquet/README.md
index 891fd02d837..b4761855f5c 100644
--- a/be/benchmark/parquet/README.md
+++ b/be/benchmark/parquet/README.md
@@ -137,6 +137,18 @@ be/output/lib/benchmark_test \
--benchmark_min_time=1s
```
+The multi-column OR pair scans the same ColumnIndex/OffsetIndex fixture and
changes only the Doris
+Page Index switch. Both variants retain the full residual expression, so the
comparison measures
+metadata pruning without changing result semantics:
+
+```shell
+be/output/lib/benchmark_test \
+ --benchmark_filter='^ParquetReader/multi_column_or/page_index_(off|on)$' \
+ --benchmark_min_time=1s \
+ --benchmark_repetitions=10 \
+ --benchmark_report_aggregates_only=true
+```
+
Every result reports throughput plus `raw_rows`, `selected_rows`,
`fixture_bytes`, `ns/raw_row`,
and (when at least one row survives) `ns/selected_row`. Keep CPU frequency,
build type, compiler,
machine placement, and benchmark filters fixed when comparing two commits.
diff --git a/be/benchmark/parquet/benchmark_parquet_reader.hpp
b/be/benchmark/parquet/benchmark_parquet_reader.hpp
index b456244ba0e..4763ffe4c98 100644
--- a/be/benchmark/parquet/benchmark_parquet_reader.hpp
+++ b/be/benchmark/parquet/benchmark_parquet_reader.hpp
@@ -34,6 +34,7 @@
#include <string>
#include <vector>
+#include "common/config.h"
#include "core/assert_cast.h"
#include "core/block/block.h"
#include "core/column/column_nullable.h"
@@ -61,6 +62,19 @@ namespace reader_detail {
constexpr size_t READER_ROWS = 1UL << 14;
constexpr size_t READER_ROW_GROUP_ROWS = 1UL << 12;
+constexpr size_t MULTI_COLUMN_OR_ROWS = 1UL << 20;
+constexpr size_t MULTI_COLUMN_OR_ROW_GROUP_ROWS = 1UL << 18;
+
+class ScopedPageIndexConfig {
+public:
+ explicit ScopedPageIndexConfig(bool enabled) :
_previous(config::enable_parquet_page_index) {
+ config::enable_parquet_page_index = enabled;
+ }
+ ~ScopedPageIndexConfig() { config::enable_parquet_page_index = _previous; }
+
+private:
+ bool _previous;
+};
inline void throw_if_error(const Status& status) {
if (!status.ok()) {
@@ -256,6 +270,8 @@ public:
_column_id(column_id),
_upper_bound(upper_bound) {}
+ bool is_constant() const override { return false; }
+
Status execute_column_impl(VExprContext*, const Block* block, const
Selector* selector,
size_t count, ColumnPtr& result_column) const
override {
DORIS_CHECK(block != nullptr);
@@ -612,6 +628,151 @@ inline void run_reader(benchmark::State& state,
ReaderScenario scenario) {
}
}
+inline std::filesystem::path ensure_multi_column_or_fixture() {
+ static std::mutex fixture_mutex;
+ const auto directory =
+ std::filesystem::temp_directory_path() /
"doris_parquet_reader_benchmark";
+ const auto path = directory / "v2_multi_column_or_page_index_v2.parquet";
+ std::lock_guard guard(fixture_mutex);
+ if (std::filesystem::exists(path)) {
+ return path;
+ }
+
+ arrow::Int32Builder ascending_builder;
+ arrow::Int32Builder descending_builder;
+ arrow::Int32Builder payload_builder;
+ PARQUET_THROW_NOT_OK(ascending_builder.Reserve(MULTI_COLUMN_OR_ROWS));
+ PARQUET_THROW_NOT_OK(descending_builder.Reserve(MULTI_COLUMN_OR_ROWS));
+ PARQUET_THROW_NOT_OK(payload_builder.Reserve(MULTI_COLUMN_OR_ROWS));
+ for (size_t row = 0; row < MULTI_COLUMN_OR_ROWS; ++row) {
+ const int32_t row_in_group = static_cast<int32_t>(row %
MULTI_COLUMN_OR_ROW_GROUP_ROWS);
+ PARQUET_THROW_NOT_OK(ascending_builder.Append(row_in_group));
+ PARQUET_THROW_NOT_OK(descending_builder.Append(
+ static_cast<int32_t>(MULTI_COLUMN_OR_ROW_GROUP_ROWS - 1) -
row_in_group));
+
PARQUET_THROW_NOT_OK(payload_builder.Append(static_cast<int32_t>(row)));
+ }
+ auto table = arrow::Table::Make(
+ arrow::schema({arrow::field("ascending", arrow::int32(), true),
+ arrow::field("descending", arrow::int32(), true),
+ arrow::field("payload", arrow::int32(), true)}),
+ {ascending_builder.Finish().ValueOrDie(),
descending_builder.Finish().ValueOrDie(),
+ payload_builder.Finish().ValueOrDie()});
+
+ std::filesystem::create_directories(directory);
+ const auto temporary_path = path.string() + ".tmp";
+ std::filesystem::remove(temporary_path);
+ const auto output_result =
arrow::io::FileOutputStream::Open(temporary_path);
+ if (!output_result.ok()) {
+ throw std::runtime_error(output_result.status().ToString());
+ }
+ const auto output = *output_result;
+ ::parquet::WriterProperties::Builder properties;
+ properties.version(::parquet::ParquetVersion::PARQUET_2_6);
+ properties.data_page_version(::parquet::ParquetDataPageVersion::V2);
+ properties.compression(::parquet::Compression::UNCOMPRESSED);
+ properties.disable_dictionary();
+ properties.encoding(::parquet::Encoding::PLAIN);
+ properties.enable_write_page_index();
+ properties.write_batch_size(8192);
+ properties.data_pagesize(64 * 1024);
+ PARQUET_THROW_NOT_OK(::parquet::arrow::WriteTable(*table,
arrow::default_memory_pool(), output,
+
MULTI_COLUMN_OR_ROW_GROUP_ROWS,
+ properties.build()));
+ PARQUET_THROW_NOT_OK(output->Close());
+ std::filesystem::rename(temporary_path, path);
+ return path;
+}
+
+inline std::unique_ptr<ReaderSession> open_multi_column_or_reader(
+ const std::filesystem::path& path) {
+ auto session = std::make_unique<ReaderSession>();
+ auto properties = std::make_shared<io::FileSystemProperties>();
+ properties->system_type = TFileType::FILE_LOCAL;
+ auto description = std::make_unique<io::FileDescription>();
+ description->path = path.string();
+ description->file_size =
static_cast<int64_t>(std::filesystem::file_size(path));
+ description->range_start_offset = 0;
+ description->range_size = -1;
+ session->reader =
std::make_unique<format::parquet::ParquetReader>(properties, description,
+
nullptr, nullptr);
+ throw_if_error(session->reader->init(&session->runtime_state));
+ throw_if_error(session->reader->get_schema(&session->schema));
+
+ session->request = std::make_shared<format::FileScanRequest>();
+ format::FileScanRequestBuilder request_builder(session->request.get());
+ std::array<int, 2> predicate_positions {};
+ for (int column = 0; column < 2; ++column) {
+ const auto column_id = format::LocalColumnId(column);
+ throw_if_error(request_builder.add_predicate_column(column_id));
+ session->request->predicate_only_columns.push_back(column_id);
+ predicate_positions[column] =
+
static_cast<int>(session->request->local_positions.at(column_id).value());
+ }
+
throw_if_error(request_builder.add_non_predicate_column(format::LocalColumnId(2)));
+
+ TExprNode node;
+ node.__set_node_type(TExprNodeType::COMPOUND_PRED);
+ node.__set_opcode(TExprOpcode::COMPOUND_OR);
+ node.__set_type(std::make_shared<DataTypeUInt8>()->to_thrift());
+ node.__set_num_children(2);
+ node.__set_is_nullable(false);
+ auto compound = VCompoundPred::create_shared(node);
+ constexpr int32_t UPPER_BOUND =
static_cast<int32_t>(MULTI_COLUMN_OR_ROW_GROUP_ROWS / 10);
+
compound->add_child(std::make_shared<Int32LessThanExpr>(predicate_positions[0],
UPPER_BOUND));
+
compound->add_child(std::make_shared<Int32LessThanExpr>(predicate_positions[1],
UPPER_BOUND));
+ auto context = VExprContext::create_shared(std::move(compound));
+ throw_if_error(context->prepare(&session->runtime_state, RowDescriptor()));
+ throw_if_error(context->open(&session->runtime_state));
+ session->request->conjuncts.push_back(context);
+ session->opened_conjuncts.push_back(std::move(context));
+ throw_if_error(session->reader->open(session->request));
+ return session;
+}
+
+inline void run_multi_column_or_reader(benchmark::State& state, bool
enable_page_index) {
+ try {
+ const auto fixture = ensure_multi_column_or_fixture();
+ ScopedPageIndexConfig page_index_config(enable_page_index);
+ size_t selected_rows = 0;
+ for (auto _ : state) {
+ state.PauseTiming();
+ auto session = open_multi_column_or_reader(fixture);
+ state.ResumeTiming();
+ const ReaderScenario scenario {.operation =
ReaderOperation::PREDICATE_SCAN,
+ .encoding = Encoding::PLAIN,
+ .null_percent = 0,
+ .null_pattern = Pattern::CLUSTERED,
+ .selectivity_percent = 20,
+ .projection =
Projection::PREDICATE_ONLY,
+ .schema_width = 3,
+ .predicate_position = 0};
+ selected_rows = scan_reader(session.get(), scenario);
+ state.PauseTiming();
+ throw_if_error(session->reader->close());
+ state.ResumeTiming();
+ benchmark::ClobberMemory();
+ }
+ constexpr size_t ROW_GROUPS = MULTI_COLUMN_OR_ROWS /
MULTI_COLUMN_OR_ROW_GROUP_ROWS;
+ constexpr size_t EXPECTED_ROWS = ROW_GROUPS * 2 *
(MULTI_COLUMN_OR_ROW_GROUP_ROWS / 10);
+ if (selected_rows != EXPECTED_ROWS) {
+ state.SkipWithError("multi-column OR benchmark returned unexpected
rows");
+ return;
+ }
+ state.SetItemsProcessed(static_cast<int64_t>(state.iterations() *
selected_rows));
+ state.counters["raw_rows"] = static_cast<double>(MULTI_COLUMN_OR_ROWS);
+ state.counters["selected_rows"] = static_cast<double>(selected_rows);
+ state.counters["fixture_bytes"] =
static_cast<double>(std::filesystem::file_size(fixture));
+ state.counters["ns/raw_row"] = benchmark::Counter(
+ static_cast<double>(MULTI_COLUMN_OR_ROWS),
+ benchmark::Counter::kIsIterationInvariantRate |
benchmark::Counter::kInvert);
+ state.counters["ns/selected_row"] = benchmark::Counter(
+ static_cast<double>(selected_rows),
+ benchmark::Counter::kIsIterationInvariantRate |
benchmark::Counter::kInvert);
+ } catch (const std::exception& error) {
+ state.SkipWithError(error.what());
+ }
+}
+
inline bool register_reader_benchmarks() {
for (const auto& scenario : reader_scenarios()) {
std::string name = "ParquetReader/" + reader_scenario_name(scenario);
@@ -619,6 +780,14 @@ inline bool register_reader_benchmarks() {
run_reader(state, scenario);
})->Unit(benchmark::kNanosecond);
}
+ benchmark::RegisterBenchmark(
+ "ParquetReader/multi_column_or/page_index_off",
+ [](benchmark::State& state) { run_multi_column_or_reader(state,
false); })
+ ->Unit(benchmark::kNanosecond);
+ benchmark::RegisterBenchmark(
+ "ParquetReader/multi_column_or/page_index_on",
+ [](benchmark::State& state) { run_multi_column_or_reader(state,
true); })
+ ->Unit(benchmark::kNanosecond);
return true;
}
diff --git a/be/src/format_v2/parquet/parquet_statistics.cpp
b/be/src/format_v2/parquet/parquet_statistics.cpp
index 75d159d670d..63b880f7c5e 100644
--- a/be/src/format_v2/parquet/parquet_statistics.cpp
+++ b/be/src/format_v2/parquet/parquet_statistics.cpp
@@ -1370,6 +1370,38 @@ std::vector<RowRange> intersect_ranges(const
std::vector<RowRange>& left,
return result;
}
+std::vector<RowRange> union_ranges(const std::vector<RowRange>& left,
+ const std::vector<RowRange>& right) {
+ std::vector<RowRange> result;
+ result.reserve(left.size() + right.size());
+ auto append = [&](const RowRange& range) {
+ if (range.length == 0) {
+ return;
+ }
+ if (!result.empty()) {
+ auto& previous = result.back();
+ const int64_t previous_end = previous.start + previous.length;
+ if (range.start <= previous_end) {
+ previous.length =
+ std::max(previous_end, range.start + range.length) -
previous.start;
+ return;
+ }
+ }
+ result.push_back(range);
+ };
+ size_t left_idx = 0;
+ size_t right_idx = 0;
+ while (left_idx < left.size() || right_idx < right.size()) {
+ if (right_idx == right.size() ||
+ (left_idx < left.size() && left[left_idx].start <=
right[right_idx].start)) {
+ append(left[left_idx++]);
+ } else {
+ append(right[right_idx++]);
+ }
+ }
+ return result;
+}
+
int64_t count_range_rows(const std::vector<RowRange>& ranges) {
int64_t rows = 0;
for (const auto& range : ranges) {
@@ -1626,6 +1658,154 @@ RowRange native_page_row_range(const
tparquet::OffsetIndex& offset_index, size_t
return {.start = start, .length = end - start};
}
+class NativePageIndexPredicateEvaluator {
+public:
+ NativePageIndexPredicateEvaluator(
+ const tparquet::FileMetaData& metadata,
+ const std::unordered_map<int, NativeParquetPageIndex>&
page_indexes,
+ const std::vector<std::unique_ptr<ParquetColumnSchema>>&
file_schema,
+ const format::FileScanRequest& request, int64_t row_group_rows,
+ ParquetPruningStats* pruning_stats, const cctz::time_zone*
timezone)
+ : _metadata(metadata),
+ _page_indexes(page_indexes),
+ _file_schema(file_schema),
+ _request(request),
+ _row_group_rows(row_group_rows),
+ _pruning_stats(pruning_stats),
+ _timezone(timezone) {}
+
+ std::optional<std::vector<RowRange>> evaluate(const VExprSPtr& expr) const
{
+ if (expr == nullptr || !expr->can_evaluate_zonemap_filter()) {
+ return std::nullopt;
+ }
+ if (expr->op() == TExprOpcode::COMPOUND_AND) {
+ return evaluate_compound(expr, true);
+ }
+ if (expr->op() == TExprOpcode::COMPOUND_OR) {
+ return evaluate_compound(expr, false);
+ }
+ return evaluate_leaf(expr);
+ }
+
+private:
+ struct SlotPageZoneMaps {
+ DataTypePtr data_type;
+ std::vector<RowRange> ranges;
+ std::vector<std::shared_ptr<segment_v2::ZoneMap>> zone_maps;
+ };
+
+ std::optional<std::vector<RowRange>> evaluate_compound(const VExprSPtr&
expr,
+ bool is_and) const {
+ std::optional<std::vector<RowRange>> ranges;
+ for (const auto& child : expr->children()) {
+ if (!child->can_evaluate_zonemap_filter()) {
+ if (!is_and) {
+ return std::nullopt;
+ }
+ continue;
+ }
+ auto child_ranges = evaluate(child);
+ if (!child_ranges.has_value()) {
+ // An unavailable AND child can be ignored, while an
unavailable OR branch must
+ // retain the complete range so metadata pruning cannot create
a false negative.
+ if (!is_and) {
+ return std::nullopt;
+ }
+ continue;
+ }
+ if (!ranges.has_value()) {
+ ranges = std::move(*child_ranges);
+ } else if (is_and) {
+ ranges = intersect_ranges(*ranges, *child_ranges);
+ } else {
+ ranges = union_ranges(*ranges, *child_ranges);
+ }
+ if (is_and && ranges->empty()) {
+ return ranges;
+ }
+ }
+ return ranges;
+ }
+
+ std::optional<std::vector<RowRange>> evaluate_leaf(const VExprSPtr& expr)
const {
+ std::set<int> slot_indexes;
+ expr->collect_slot_column_ids(slot_indexes);
+ if (slot_indexes.size() != 1) {
+ return std::nullopt;
+ }
+ const int slot_index = *slot_indexes.begin();
+ const auto* pages = load_slot_pages(slot_index);
+ if (pages == nullptr) {
+ return std::nullopt;
+ }
+
+ std::vector<RowRange> ranges;
+ for (size_t page_idx = 0; page_idx < pages->ranges.size(); ++page_idx)
{
+ ZoneMapEvalContext ctx;
+ add_slot_zonemap(&ctx, slot_index, pages->data_type,
pages->zone_maps[page_idx]);
+ if (expr->evaluate_zonemap_filter(ctx) !=
ZoneMapFilterResult::kNoMatch) {
+ append_row_range(pages->ranges[page_idx], &ranges);
+ }
+ accumulate_zonemap_stats(ctx, _pruning_stats);
+ }
+ return ranges;
+ }
+
+ const SlotPageZoneMaps* load_slot_pages(int slot_index) const {
+ const auto cached = _slot_page_zone_maps.find(slot_index);
+ if (cached != _slot_page_zone_maps.end()) {
+ return cached->second.has_value() ? &*cached->second : nullptr;
+ }
+ const auto file_column_id = file_column_id_by_block_position(_request,
slot_index);
+ if (!file_column_id.has_value()) {
+ _slot_page_zone_maps.emplace(slot_index, std::nullopt);
+ return nullptr;
+ }
+ const auto* column_schema = resolve_local_leaf_schema(_file_schema,
*file_column_id);
+ if (column_schema == nullptr || column_schema->type == nullptr ||
+ !native_metadata_predicate_is_type_safe(*column_schema) ||
+ !detail::has_supported_type_defined_order(_metadata,
column_schema->leaf_column_id)) {
+ _slot_page_zone_maps.emplace(slot_index, std::nullopt);
+ return nullptr;
+ }
+ const auto index_it =
_page_indexes.find(column_schema->leaf_column_id);
+ if (index_it == _page_indexes.end()) {
+ _slot_page_zone_maps.emplace(slot_index, std::nullopt);
+ return nullptr;
+ }
+
+ const auto& indexes = index_it->second;
+ SlotPageZoneMaps pages;
+ pages.data_type = column_schema->type;
+ pages.ranges.reserve(indexes.offset_index.page_locations.size());
+ pages.zone_maps.reserve(indexes.offset_index.page_locations.size());
+ for (size_t page_idx = 0; page_idx <
indexes.offset_index.page_locations.size();
+ ++page_idx) {
+ const auto page_range =
+ native_page_row_range(indexes.offset_index, page_idx,
_row_group_rows);
+ ParquetColumnStatistics statistics;
+ if (!build_native_page_statistics(indexes.column_index,
*column_schema, page_idx,
+ page_range.length, &statistics,
_timezone)) {
+ _slot_page_zone_maps.emplace(slot_index, std::nullopt);
+ return nullptr;
+ }
+ pages.ranges.push_back(page_range);
+
pages.zone_maps.push_back(ParquetStatisticsUtils::MakeZoneMap(statistics));
+ }
+ const auto inserted = _slot_page_zone_maps.emplace(slot_index,
std::move(pages));
+ return &*inserted.first->second;
+ }
+
+ const tparquet::FileMetaData& _metadata;
+ const std::unordered_map<int, NativeParquetPageIndex>& _page_indexes;
+ const std::vector<std::unique_ptr<ParquetColumnSchema>>& _file_schema;
+ const format::FileScanRequest& _request;
+ int64_t _row_group_rows;
+ ParquetPruningStats* _pruning_stats;
+ const cctz::time_zone* _timezone;
+ mutable std::unordered_map<int, std::optional<SlotPageZoneMaps>>
_slot_page_zone_maps;
+};
+
} // namespace
Status select_row_group_ranges_by_native_page_index(
@@ -1654,10 +1834,14 @@ Status select_row_group_ranges_by_native_page_index(
}
std::map<int, VExprContextSPtrs> conjuncts_by_slot;
+ VExprContextSPtrs multi_slot_conjuncts;
for (const auto& conjunct : request.conjuncts) {
const auto slot_index =
expr_zonemap::single_slot_zonemap_index(conjunct);
if (slot_index >= 0) {
conjuncts_by_slot[slot_index].push_back(conjunct);
+ } else if (conjunct != nullptr && conjunct->root() != nullptr &&
+ conjunct->root()->can_evaluate_zonemap_filter()) {
+ multi_slot_conjuncts.push_back(conjunct);
}
}
for (const auto& [slot_index, conjuncts] : conjuncts_by_slot) {
@@ -1695,12 +1879,7 @@ Status select_row_group_ranges_by_native_page_index(
ZoneMapFilterResult::kNoMatch) {
append_row_range(page_range, &filter_ranges);
}
- if (pruning_stats != nullptr) {
- pruning_stats->expr_zonemap_unusable_evals +=
ctx.stats.unusable_zonemap_eval_count;
- pruning_stats->in_zonemap_point_check_count +=
- ctx.stats.in_zonemap_point_check_count;
- pruning_stats->in_zonemap_range_only_count +=
ctx.stats.in_zonemap_range_only_count;
- }
+ accumulate_zonemap_stats(ctx, pruning_stats);
}
if (!usable) {
continue;
@@ -1715,6 +1894,23 @@ Status select_row_group_ranges_by_native_page_index(
}
}
+ NativePageIndexPredicateEvaluator evaluator(metadata, page_indexes,
file_schema, request,
+ row_group_rows, pruning_stats,
timezone);
+ for (const auto& conjunct : multi_slot_conjuncts) {
+ auto conjunct_ranges = evaluator.evaluate(conjunct->root());
+ if (!conjunct_ranges.has_value()) {
+ continue;
+ }
+ *selected_ranges = intersect_ranges(*selected_ranges,
*conjunct_ranges);
+ if (selected_ranges->empty()) {
+ if (pruning_stats != nullptr) {
+ pruning_stats->filtered_page_rows += row_group_rows;
+ ++pruning_stats->filtered_row_groups_by_page_index;
+ }
+ return Status::OK();
+ }
+ }
+
for (const auto& conjunct : request.conjuncts) {
const auto predicate = extract_variant_shredded_predicate(conjunct);
if (!predicate.has_value()) {
diff --git a/be/test/format_v2/parquet/parquet_scan_test.cpp
b/be/test/format_v2/parquet/parquet_scan_test.cpp
index 8ed63a7fe28..c3532c051f0 100644
--- a/be/test/format_v2/parquet/parquet_scan_test.cpp
+++ b/be/test/format_v2/parquet/parquet_scan_test.cpp
@@ -1678,6 +1678,17 @@ void write_page_index_parquet_file(const std::string&
file_path) {
write_table(file_path, table, ids.size(), false, true);
}
+void write_multi_column_page_index_parquet_file(const std::string& file_path) {
+ std::vector<int32_t> ascending(128);
+ std::iota(ascending.begin(), ascending.end(), 0);
+ std::vector<int32_t> descending(ascending.rbegin(), ascending.rend());
+ auto schema = arrow::schema({arrow::field("ascending", arrow::int32(),
false),
+ arrow::field("descending", arrow::int32(),
false)});
+ auto table = arrow::Table::Make(schema,
+ {build_int32_array(ascending),
build_int32_array(descending)});
+ write_table(file_path, table, ascending.size(), false, true);
+}
+
void write_multi_row_group_page_index_parquet_file(const std::string&
file_path) {
std::vector<int32_t> ids(384);
std::iota(ids.begin(), ids.end(), 0);
@@ -4193,6 +4204,52 @@ TEST_F(ParquetScanTest,
ProfileCountersReflectPageIndexAndRangeGapPruning) {
EXPECT_GT(profile.get_counter("RangeGapSkippedRows")->value(), 0);
}
+TEST_F(ParquetScanTest, MultiColumnOrUsesPageIndexAndResidualExpression) {
+ write_multi_column_page_index_parquet_file(_file_path);
+ RuntimeProfile profile("profile");
+ auto reader = create_reader(0, -1, &profile);
+ RuntimeState state {TQueryOptions(), TQueryGlobals()};
+ ASSERT_TRUE(reader->init(&state).ok());
+
+ std::vector<format::ColumnDefinition> schema;
+ ASSERT_TRUE(reader->get_schema(&schema).ok());
+ auto request = std::make_shared<format::FileScanRequest>();
+ format::FileScanRequestBuilder request_builder(request.get());
+
ASSERT_TRUE(request_builder.add_predicate_column(format::LocalColumnId(0)).ok());
+
ASSERT_TRUE(request_builder.add_predicate_column(format::LocalColumnId(1)).ok());
+ const auto low_ascending = create_int32_function_conjunct(0, "lt",
TExprOpcode::LT, 13);
+ const auto low_descending = create_int32_function_conjunct(1, "lt",
TExprOpcode::LT, 13);
+ auto disjunction = create_compound_conjunct(TExprOpcode::COMPOUND_OR,
low_ascending->root(),
+ low_descending->root());
+ ASSERT_TRUE(disjunction->prepare(&state, RowDescriptor()).ok());
+ ASSERT_TRUE(disjunction->open(&state).ok());
+ request->conjuncts.push_back(disjunction);
+ ASSERT_TRUE(reader->open(request).ok());
+
+ std::vector<int32_t> selected;
+ bool eof = false;
+ while (!eof) {
+ Block block = build_file_block(schema);
+ size_t rows = 0;
+ ASSERT_TRUE(reader->get_block(&block, &rows, &eof).ok());
+ const auto& values =
int32_data_column(*block.get_by_position(0).column);
+ for (size_t row = 0; row < rows; ++row) {
+ selected.push_back(values.get_element(row));
+ }
+ }
+
+ std::vector<int32_t> expected(13);
+ std::iota(expected.begin(), expected.end(), 0);
+ for (int32_t value = 115; value < 128; ++value) {
+ expected.push_back(value);
+ }
+ EXPECT_EQ(selected, expected);
+ EXPECT_GT(counter_value(profile, "FilteredRowsByPage"), 0);
+ EXPECT_LT(counter_value(profile, "RawRowsRead"), 128);
+ EXPECT_GT(counter_value(profile, "RowsFilteredByConjunct"), 0);
+ disjunction->close();
+}
+
TEST_F(ParquetScanTest, OpenDefersPageIndexProbeToCurrentRowGroup) {
write_multi_row_group_page_index_parquet_file(_file_path);
RuntimeProfile profile("lazy_page_index_profile");
diff --git a/be/test/format_v2/parquet/parquet_statistics_test.cpp
b/be/test/format_v2/parquet/parquet_statistics_test.cpp
index 4a6dc62775e..8f7e8c2b285 100644
--- a/be/test/format_v2/parquet/parquet_statistics_test.cpp
+++ b/be/test/format_v2/parquet/parquet_statistics_test.cpp
@@ -40,6 +40,7 @@
#include "core/data_type/data_type_variant_v2.h"
#include "core/field.h"
#include "exprs/expr_zonemap_filter.h"
+#include "exprs/vcompound_pred.h"
#include "exprs/vexpr.h"
#include "exprs/vexpr_context.h"
#include "exprs/vliteral.h"
@@ -167,6 +168,41 @@ private:
const std::string _expr_name = "MetadataInt32GreaterThanExpr";
};
+class MetadataSlotInt32GreaterThanExpr final : public VExpr {
+public:
+ MetadataSlotInt32GreaterThanExpr(int slot_index, int32_t value)
+ : VExpr(std::make_shared<DataTypeUInt8>(), false),
+ _slot_index(slot_index),
+ _value(value) {}
+
+ const std::string& expr_name() const override { return _expr_name; }
+ Status execute_column_impl(VExprContext*, const Block*, const Selector*,
size_t,
+ ColumnPtr&) const override {
+ return Status::InternalError("MetadataSlotInt32GreaterThanExpr is
metadata-only");
+ }
+ bool can_evaluate_zonemap_filter() const override { return true; }
+ void collect_slot_column_ids(std::set<int>& column_ids) const override {
+ column_ids.insert(_slot_index);
+ }
+ ZoneMapFilterResult evaluate_zonemap_filter(const ZoneMapEvalContext& ctx)
const override {
+ const auto zone_map = ctx.zone_map(_slot_index);
+ if (zone_map == nullptr) {
+ return unsupported_zonemap_filter(ctx);
+ }
+ if (!zone_map->has_not_null) {
+ return ZoneMapFilterResult::kNoMatch;
+ }
+ return zone_map->max_value <= Field::create_field<TYPE_INT>(_value)
+ ? ZoneMapFilterResult::kNoMatch
+ : ZoneMapFilterResult::kMayMatch;
+ }
+
+private:
+ int _slot_index;
+ int32_t _value;
+ const std::string _expr_name = "MetadataSlotInt32GreaterThanExpr";
+};
+
class MetadataBoundsProbeExpr final : public VExpr {
public:
explicit MetadataBoundsProbeExpr(bool require_false_boolean = false)
@@ -466,6 +502,130 @@ TEST(NativeParquetStatisticsTest,
InvalidTimeAndPaddedBooleanPageBoundsCannotPru
EXPECT_EQ(bool_ranges[0].start, 0);
EXPECT_EQ(bool_ranges[0].length, 1);
}
+
+TEST(NativeParquetStatisticsTest, MultiColumnOrUnionsPageIndexRanges) {
+ auto encode_int32 = [](int32_t value) {
+ std::string bytes(sizeof(value), '\0');
+ memcpy(bytes.data(), &value, sizeof(value));
+ return bytes;
+ };
+ auto make_schema = [](int local_id, int leaf_column_id) {
+ auto column = std::make_unique<format::parquet::ParquetColumnSchema>();
+ column->kind = format::parquet::ParquetColumnSchemaKind::PRIMITIVE;
+ column->local_id = local_id;
+ column->leaf_column_id = leaf_column_id;
+ column->type = std::make_shared<DataTypeInt32>();
+ column->type_descriptor.doris_type = column->type;
+ column->type_descriptor.physical_type = tparquet::Type::INT32;
+ return column;
+ };
+ auto make_page_index = [&](const std::vector<int32_t>& values) {
+ format::parquet::NativeParquetPageIndex page_index;
+ std::vector<std::string> encoded;
+ encoded.reserve(values.size());
+ for (const auto value : values) {
+ encoded.push_back(encode_int32(value));
+ }
+ page_index.column_index.__set_min_values(encoded);
+ page_index.column_index.__set_max_values(encoded);
+
page_index.column_index.__set_null_pages(std::vector<bool>(values.size(),
false));
+
page_index.column_index.__set_null_counts(std::vector<int64_t>(values.size(),
0));
+ std::vector<tparquet::PageLocation> locations;
+ for (size_t page_idx = 0; page_idx < values.size(); ++page_idx) {
+ tparquet::PageLocation location;
+ location.__set_offset(static_cast<int64_t>(page_idx * 100));
+ location.__set_compressed_page_size(100);
+ location.__set_first_row_index(static_cast<int64_t>(page_idx *
10));
+ locations.push_back(location);
+ }
+ page_index.offset_index.__set_page_locations(std::move(locations));
+ return page_index;
+ };
+
+ std::vector<std::unique_ptr<format::parquet::ParquetColumnSchema>> schema;
+ schema.push_back(make_schema(0, 0));
+ schema.push_back(make_schema(1, 1));
+ tparquet::ColumnOrder order;
+ order.__set_TYPE_ORDER(tparquet::TypeDefinedOrder());
+ tparquet::FileMetaData metadata;
+ metadata.__set_column_orders({order, order});
+
+ auto make_compound_expr = [](TExprOpcode::type opcode, VExprSPtr left,
VExprSPtr right) {
+ TExprNode compound_node;
+ compound_node.__set_node_type(TExprNodeType::COMPOUND_PRED);
+ compound_node.__set_opcode(opcode);
+
compound_node.__set_type(std::make_shared<DataTypeUInt8>()->to_thrift());
+ compound_node.__set_num_children(2);
+ compound_node.__set_is_nullable(false);
+ auto compound = VCompoundPred::create_shared(compound_node);
+ compound->add_child(std::move(left));
+ compound->add_child(std::move(right));
+ return compound;
+ };
+ auto make_compound = [&](TExprOpcode::type opcode) {
+ return VExprContext::create_shared(make_compound_expr(
+ opcode, std::make_shared<MetadataSlotInt32GreaterThanExpr>(0,
50),
+ std::make_shared<MetadataSlotInt32GreaterThanExpr>(1, 50)));
+ };
+
+ format::FileScanRequest request;
+ request.local_positions.emplace(format::LocalColumnId(0),
format::LocalIndex(0));
+ request.local_positions.emplace(format::LocalColumnId(1),
format::LocalIndex(1));
+ request.predicate_columns =
{format::LocalColumnIndex::top_level(format::LocalColumnId(0)),
+
format::LocalColumnIndex::top_level(format::LocalColumnId(1))};
+ request.conjuncts = {make_compound(TExprOpcode::COMPOUND_OR)};
+
+ std::unordered_map<int, format::parquet::NativeParquetPageIndex>
page_indexes;
+ page_indexes.emplace(0, make_page_index({100, 0, 0}));
+ page_indexes.emplace(1, make_page_index({0, 0, 100}));
+ std::vector<format::parquet::RowRange> selected_ranges;
+ std::map<int, format::parquet::ParquetPageSkipPlan> skip_plans;
+ ASSERT_TRUE(format::parquet::select_row_group_ranges_by_native_page_index(
+ metadata, tparquet::RowGroup {}, page_indexes, schema,
request, 30,
+ &selected_ranges, &skip_plans, nullptr)
+ .ok());
+ ASSERT_EQ(selected_ranges.size(), 2);
+ EXPECT_EQ(selected_ranges[0].start, 0);
+ EXPECT_EQ(selected_ranges[0].length, 10);
+ EXPECT_EQ(selected_ranges[1].start, 20);
+ EXPECT_EQ(selected_ranges[1].length, 10);
+
+ page_indexes.erase(1);
+ ASSERT_TRUE(format::parquet::select_row_group_ranges_by_native_page_index(
+ metadata, tparquet::RowGroup {}, page_indexes, schema,
request, 30,
+ &selected_ranges, &skip_plans, nullptr)
+ .ok());
+ ASSERT_EQ(selected_ranges.size(), 1);
+ EXPECT_EQ(selected_ranges[0].start, 0);
+ EXPECT_EQ(selected_ranges[0].length, 30);
+
+ request.conjuncts = {make_compound(TExprOpcode::COMPOUND_AND)};
+ ASSERT_TRUE(format::parquet::select_row_group_ranges_by_native_page_index(
+ metadata, tparquet::RowGroup {}, page_indexes, schema,
request, 30,
+ &selected_ranges, &skip_plans, nullptr)
+ .ok());
+ ASSERT_EQ(selected_ranges.size(), 1);
+ EXPECT_EQ(selected_ranges[0].start, 0);
+ EXPECT_EQ(selected_ranges[0].length, 10);
+
+ page_indexes.emplace(1, make_page_index({0, 0, 100}));
+ auto first_branch = make_compound_expr(
+ TExprOpcode::COMPOUND_AND,
std::make_shared<MetadataSlotInt32GreaterThanExpr>(0, 50),
+ std::make_shared<MetadataSlotInt32GreaterThanExpr>(1, 50));
+ auto second_branch = make_compound_expr(
+ TExprOpcode::COMPOUND_AND,
std::make_shared<MetadataSlotInt32GreaterThanExpr>(0, -1),
+ std::make_shared<MetadataSlotInt32GreaterThanExpr>(1, 50));
+ request.conjuncts = {VExprContext::create_shared(make_compound_expr(
+ TExprOpcode::COMPOUND_OR, std::move(first_branch),
std::move(second_branch)))};
+ ASSERT_TRUE(format::parquet::select_row_group_ranges_by_native_page_index(
+ metadata, tparquet::RowGroup {}, page_indexes, schema,
request, 30,
+ &selected_ranges, &skip_plans, nullptr)
+ .ok());
+ ASSERT_EQ(selected_ranges.size(), 1);
+ EXPECT_EQ(selected_ranges[0].start, 20);
+ EXPECT_EQ(selected_ranges[0].length, 10);
+}
+
TEST(ParquetBloomFilterPruningTest, NativeUint32BloomUsesPhysicalInt32Hash) {
const auto column_schema = uint32_parquet_bloom_schema();
format::parquet::native::BlockSplitBloomFilter bloom_filter;
diff --git a/docs/file-scanner-v2-parquet-scan-design.md
b/docs/file-scanner-v2-parquet-scan-design.md
index b81a779ac5e..29dfe66c6fa 100644
--- a/docs/file-scanner-v2-parquet-scan-design.md
+++ b/docs/file-scanner-v2-parquet-scan-design.md
@@ -197,9 +197,9 @@ flowchart LR
does not repeatedly interpret table-schema evolution.
3. **Capability checks:** ZoneMap, Dictionary, and Bloom use only expressions
they can interpret
safely. All others remain row-level residual predicates.
-4. **Prefer safe single-column predicates:** Single-column predicates can
drive indexes and staged
- filtering. Multi-column, stateful, or error-sensitive expressions retain
whole-expression
- evaluation.
+4. **Prefer safe single-column row filters:** Single-column predicates can
drive staged dictionary
+ or raw filtering. Multi-column AND/OR trees may still combine conservative
Row Group and Page
+ Index candidate ranges, but the complete expression remains in
whole-expression row evaluation.
5. **Runtime Filters can refresh:** ScannerScheduler refreshes late Runtime
Filters before reading.
TableReader handles partition-range pruning during Split preparation, and
passes file-pushable
parts as localized conjuncts.
@@ -285,9 +285,11 @@ sequenceDiagram
### How the plan drives physical skips
ColumnIndex provides min/max/null semantics for each page. OffsetIndex maps
pages to Row Group row
-numbers and file offsets. Candidate ranges from multiple predicate columns are
intersected into
-`selected_ranges`; a `page_skip_plan` is then built for each leaf so its
column reader can skip pages
-that do not overlap surviving rows.
+numbers and file offsets. Candidate ranges follow the predicate tree: AND
nodes intersect child
+ranges and OR nodes union them into `selected_ranges`. A missing or unusable
AND child contributes
+no pruning, while a missing or unusable OR branch retains the complete Row
Group range. A
+`page_skip_plan` is then built for each leaf so its column reader can skip
pages that do not overlap
+surviving rows.
> `selected_ranges` represents logical row ranges, while `page_skip_plan`
> represents physical page
> reads. Keeping them separate allows the scheduler to advance by row batch
> while each column skips
@@ -826,7 +828,8 @@ split safely, or read anomalies must never change query
semantics.
| Bloom missing, disabled, or unreadable | Skip Bloom pruning and continue
with later scan stages |
| Incomplete dictionary page, mixed non-dictionary encoding, complex/repeated
column | Disable dictionary pruning and Dictionary-ID Filter; use actual values
|
| Missing or inconsistent ColumnIndex/OffsetIndex | Disable fine-grained page
pruning and read the full candidate range |
-| Multi-column, OR, stateful, or error-order-sensitive expression | Preserve
whole-expression evaluation to avoid changing SQL short-circuit or error
semantics |
+| Multi-column AND/OR expression | Combine only conservative metadata
candidate ranges; preserve whole-expression row evaluation |
+| Stateful or error-order-sensitive expression | Preserve whole-expression
evaluation without metadata decomposition |
| No stable file-version identity for Page Cache | Disable Parquet Page Cache
to prevent stale-byte reads |
| Incomplete Condition Cache coverage | Retain and recompute uncovered ranges |
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]