Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
65399e1
[improvement](hive) Support partition-column-value-only pushdown for …
Sep 23, 2026
fb90f6a
limit PhysicalPartitionTopN runtime filter effective scope
Sep 23, 2026
1a4bf1c
fix code style
Sep 30, 2026
fcc73b2
[fix](hive) Require proven nonempty ranges for partition-column-value…
Oct 1, 2026
fab967b
fix code style
Oct 1, 2026
c5d6cf7
[fix](regression) Split Hive SET from INSERT in partition-value fixture
Oct 1, 2026
0aa465d
[fix](hive) Keep ordinary file splitting for partition-value scans
Oct 1, 2026
b13b9fb
[fix](nereids) Inherit CTE effectiveness only from a bounded producer
Oct 1, 2026
645ec9d
fix be ut
Oct 2, 2026
070bf4c
fix code style
Oct 2, 2026
ea7a6de
[fix](nereids) Propagate CTE effectiveness from relation-global produ…
Oct 2, 2026
f92db8a
[fix](nereids) Carry a proven bound through grouped aggregates and wi…
Oct 3, 2026
94e4804
[test](hive) Assert the partition-value pushdown actually ran
Oct 3, 2026
95d70de
[test](hive) Cover the partition-value pushdown on scanner V1
Oct 3, 2026
15c91ea
[style] Sort RuntimeFilterTest imports
Oct 3, 2026
9a4fd36
[improvement](hive) Emit a partition value per range, and cover Hudi
Oct 3, 2026
034d144
[improvement](hive) Keep the partition-value pushdown with a retained…
Oct 3, 2026
a37119c
[improvement](hive) Support COUNT(DISTINCT) in the partition-value pu…
Oct 3, 2026
a3689d1
fix code style
Oct 3, 2026
dc151f9
[test](hive) Cover the two-consumer plans from the CTE-bound review
Oct 3, 2026
aa7e75f
[fix] Update the tests that encoded the old partition-value semantics
Oct 3, 2026
a073437
fix code style
Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions be/src/exec/scan/file_scanner.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@
#include "format/json/new_json_reader.h"
#include "format/orc/vorc_reader.h"
#include "format/parquet/vparquet_reader.h"
#include "format/partition_column_reader.h"
#include "format/table/es/es_http_reader.h"
#include "format/table/hive_reader.h"
#include "format/table/hudi_jni_reader.h"
Expand Down Expand Up @@ -1310,6 +1311,24 @@ Status FileScanner::_get_next_reader() {
}
}

// A partition value is an input row for every range the partition has. This no longer needs
// the footer to prove nonemptiness, so it also no longer needs a reader that can report a
// row count; has_delete_operations() still keeps formats with row-level deletes out, and
// supports_range() keeps non Hive/Hudi formats and non Parquet/ORC files out.
if (_get_push_down_agg_type() == TPushAggOp::type::PARTITION_VALUE &&
PartitionColumnReader::supports_range(range, format_type) &&
!_partition_col_descs.empty() && _file_slot_descs.empty() && _conjuncts.empty() &&
_applied_rf_num == _total_rf_num && !_cur_reader->has_delete_operations() &&
std::all_of(_column_descs.begin(), _column_descs.end(),
[this](const ColumnDescriptor& col_desc) {
return col_desc.category == ColumnCategory::PARTITION_KEY &&
_partition_col_descs.contains(col_desc.name);
})) {
auto* table_reader = assert_cast<TableFormatReader*>(_cur_reader.release());
_cur_reader = std::make_unique<PartitionColumnReader>(
std::unique_ptr<TableFormatReader>(table_reader));
}

// Unified COUNT(*) pushdown: replace the real reader with CountReader
// decorator if the reader accepts COUNT and can provide a total row count.
if (_cur_reader->get_push_down_agg_type() == TPushAggOp::type::COUNT) {
Expand Down
69 changes: 69 additions & 0 deletions be/src/format/partition_column_reader.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

#pragma once

#include <cstddef>
#include <cstdint>
#include <memory>

#include "format/count_reader.h"
#include "format/table/table_format_reader.h"

namespace doris {

// Decorates an initialized Hive/Hudi reader so a partition-only duplicate-insensitive aggregate is
// answered from partition metadata instead of file data.
//
// One row of partition values is emitted per scan range, unconditionally. A range whose file turns
// out to hold zero rows still contributes its partition value, so MAX/GROUP BY/DISTINCT can name a
// partition that a full scan would not return. A partition with no file at all contributes nothing,
// because it produces no scan range.
class PartitionColumnReader final : public CountReader {
public:
// The reader must be initialized (footer parsed) before it can fill the typed partition values,
// but its row count is irrelevant now, so V1 no longer requires the whole file: the whole-range
// rule existed only because a partial range's count was unreliable (Parquet row groups are not
// filtered by range when counting, and ORC's count reads 0 until its row reader exists).
static bool supports_range(const TFileRangeDesc& range, TFileFormatType::type format_type) {
return range.__isset.table_format_params &&
(range.table_format_params.table_format_type == "hive" ||
range.table_format_params.table_format_type == "hudi") &&
(format_type == TFileFormatType::FORMAT_PARQUET ||
format_type == TFileFormatType::FORMAT_ORC);
}

explicit PartitionColumnReader(std::unique_ptr<TableFormatReader> inner_reader)
: CountReader(1, 1, std::move(inner_reader)) {
DORIS_CHECK(this->inner_reader() != nullptr);
set_push_down_agg_type(TPushAggOp::type::PARTITION_VALUE);
}

protected:
Status on_after_read_block(Block* block, size_t* read_rows) override {
if (*read_rows > 0) {
// CountReader supplies cardinality; the initialized reader owns typed partition values.
// Fill helpers append, so discard the default cells before materializing constants.
block->clear_column_data();
RETURN_IF_ERROR(static_cast<TableFormatReader*>(inner_reader())
->fill_remaining_columns(block, *read_rows));
}
return Status::OK();
}
};

} // namespace doris
29 changes: 29 additions & 0 deletions be/src/format_v2/table_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,8 @@ std::string push_down_agg_to_string(TPushAggOp::type op) {
return "MIX";
case TPushAggOp::COUNT_ON_INDEX:
return "COUNT_ON_INDEX";
case TPushAggOp::PARTITION_VALUE:
return "PARTITION_VALUE";
}
return "UNKNOWN";
}
Expand Down Expand Up @@ -1599,6 +1601,33 @@ Status TableReader::_evaluate_partition_prune_conjuncts(const VExprContextSPtrs&
can_filter_all);
}

bool TableReader::_conjuncts_reference_only_partition_columns() const {
for (const auto& conjunct : _conjuncts) {
if (conjunct == nullptr || conjunct->root() == nullptr) {
return false;
}
std::set<GlobalIndex> global_indices;
collect_global_indices(conjunct->root(), &global_indices);
// A slotless predicate is deliberately excluded: it would be evaluated once against the
// synthesized row instead of once per source row, which changes its row-level semantics.
if (global_indices.empty()) {
return false;
}
const bool partition_only = std::ranges::all_of(global_indices, [this](GlobalIndex index) {
if (index.value() >= _projected_columns.size()) {
return false;
}
const auto& column = _projected_columns[index.value()];
return column.is_partition_key &&
find_partition_value(column, _partition_values) != nullptr;
});
if (!partition_only) {
return false;
}
}
return true;
}

bool TableReader::_is_safe_to_pre_execute(const VExprContextSPtr& conjunct) {
DORIS_CHECK(conjunct != nullptr);
DORIS_CHECK(conjunct->root() != nullptr);
Expand Down
65 changes: 56 additions & 9 deletions be/src/format_v2/table_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -551,6 +551,10 @@ class TableReader {
Status _evaluate_partition_prune_conjuncts(const VExprContextSPtrs& conjuncts,
bool* can_filter_all);
static bool _is_safe_to_pre_execute(const VExprContextSPtr& conjunct);
// Whether every conjunct the scanner will evaluate reads nothing but partition columns.
// Only then may PARTITION_VALUE hand it the one-row block it synthesizes: see
// _supports_aggregate_pushdown(TPushAggOp::type::PARTITION_VALUE).
bool _conjuncts_reference_only_partition_columns() const;
Status _build_partition_prune_block(Block* block) const;
Status _open_local_filter_exprs(const FileScanRequest& file_request);
Status _init_reader_condition_cache(const FileScanRequest& file_request);
Expand Down Expand Up @@ -1058,6 +1062,22 @@ class TableReader {
return Status::OK();
}

// PARTITION_VALUE needs no file data at all: every projected column is a partition column,
// so one row of partition values is a faithful row of this range's output. Emitting it
// unconditionally also means the optimization no longer depends on the reader being able to
// prove a row count (Hudi and other formats do not all support count pushdown).
//
// The trade-off is deliberate: a range whose file turns out to hold zero rows still
// contributes its partition value, so MAX/GROUP BY/DISTINCT can name a partition that a full
// scan would not return. A partition with no file at all still produces nothing, because it
// produces no scan range.
if (_push_down_agg_type == TPushAggOp::type::PARTITION_VALUE) {
RETURN_IF_ERROR(finalize_chunk(block, 1));
*pushed_down = true;
RETURN_IF_ERROR(close_current_reader());
return Status::OK();
}

FileAggregateRequest file_request;
RETURN_IF_ERROR(_build_file_aggregate_request(_push_down_agg_type, &file_request));
FileAggregateResult file_result;
Expand Down Expand Up @@ -1091,8 +1111,8 @@ class TableReader {
}

virtual bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const {
// Only COUNT and MIN/MAX can be push down.
if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX) {
if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX &&
agg_type != TPushAggOp::type::PARTITION_VALUE) {
return false;
}
// Aggregate pushdown returns reduced synthetic rows and may close the physical reader
Expand All @@ -1103,19 +1123,44 @@ class TableReader {
if (!_all_runtime_filters_applied_for_split) {
return false;
}
// Scanner owns the original conjunct list and evaluates it after TableReader finalizes
// rows. Even a slotless conjunct that cannot become a TableFilter must see every source
// row before an aggregate reduces the stream to synthetic COUNT/MINMAX rows.
if (!_conjuncts.empty()) {
return false;
}
// Only support aggregate pushdown when there is no delete or filter, so
// Only support aggregate pushdown when there is no delete, so
// the reduced rows consumed by the upper aggregate remain semantically equivalent to a
// normal scan.
if ((_delete_rows != nullptr && !_delete_rows->empty()) ||
(_deletion_vector != nullptr && !_deletion_vector->isEmpty())) {
return false;
}
if (agg_type == TPushAggOp::type::PARTITION_VALUE) {
DORIS_CHECK(_file_scan_request != nullptr);
if (!_current_file_range_desc.__isset.table_format_params ||
(_current_file_range_desc.table_format_params.table_format_type != "hive" &&
_current_file_range_desc.table_format_params.table_format_type != "hudi") ||
(_format != FileFormat::PARQUET && _format != FileFormat::ORC) ||
_projected_columns.empty() || !_file_scan_request->delete_conjuncts.empty()) {
return false;
}
if (!std::ranges::all_of(_projected_columns, [this](const auto& column) {
return column.is_partition_key &&
find_partition_value(column, _partition_values) != nullptr;
})) {
return false;
}
// A retained predicate is NOT a reason to decline. Scanner::_filter_output_block()
// evaluates the scanner's conjuncts on whatever block this reader returns, so the
// one-row block synthesized for PARTITION_VALUE is filtered exactly like a real row.
// That is sound only while the conjuncts read nothing but partition columns: the
// synthesized row carries partition values and nothing else, so a predicate on a
// data column, or on a slot outside the projection, would be evaluated against
// unrelated values. Requiring at least one referenced slot also keeps a slotless
// predicate from being evaluated once here instead of once per source row.
return _conjuncts_reference_only_partition_columns();
}
// Scanner owns the original conjunct list and evaluates it after TableReader finalizes
// rows. Even a slotless conjunct that cannot become a TableFilter must see every source
// row before an aggregate reduces the stream to synthetic COUNT/MINMAX rows.
if (!_conjuncts.empty()) {
return false;
}
if (!_table_filters.empty()) {
return false;
}
Expand Down Expand Up @@ -2090,6 +2135,8 @@ class TableReader {
DORIS_CHECK(_supports_aggregate_pushdown(agg_type));
request->agg_type = agg_type;
request->columns.clear();
// PARTITION_VALUE never reaches here: _try_materialize_aggregate_pushdown_rows emits the
// partition row without asking the reader for a count.
if (agg_type == TPushAggOp::type::COUNT) {
DORIS_CHECK(_push_down_count_columns.has_value());
// An empty explicit list is the semantic signal for COUNT(*). Do not inspect the
Expand Down
93 changes: 92 additions & 1 deletion be/test/format/table/table_format_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,25 @@
#include "core/column/column_vector.h"
#include "core/data_type/data_type_nullable.h"
#include "core/data_type/data_type_number.h"
#include "format/partition_column_reader.h"
#include "runtime/descriptors.h"

namespace doris {

class MockTableFormatReader : public TableFormatReader {
public:
Status _do_get_next_block(Block*, size_t*, bool*) override { return Status::OK(); }
Status _do_get_next_block(Block*, size_t*, bool*) override {
++read_calls;
return Status::OK();
}

Status close() override {
++close_calls;
return Status::OK();
}

int read_calls = 0;
int close_calls = 0;

void set_fill_col_name_to_block_idx(std::unordered_map<std::string, uint32_t>* index) {
_fill_col_name_to_block_idx = index;
Expand Down Expand Up @@ -150,4 +162,83 @@ TEST(TableFormatReaderTest, FillMissingNullableColumnDetachesSharedBlockSlot) {
EXPECT_EQ(null_map[2], 1);
}

TEST(TableFormatReaderTest, PartitionValueSupportsColumnarHiveAndHudiRanges) {
TTableFormatFileDesc table_format;
table_format.__set_table_format_type("hive");
TFileRangeDesc range;
range.__set_table_format_params(table_format);
range.__set_start_offset(0);
range.__set_size(1024);
range.__set_file_size(1024);
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_PARQUET));
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_ORC));
// Only formats whose reader can be opened and whose metadata is readable.
for (const auto format : {TFileFormatType::FORMAT_CSV_PLAIN, TFileFormatType::FORMAT_TEXT,
TFileFormatType::FORMAT_JSON, TFileFormatType::FORMAT_JNI}) {
EXPECT_FALSE(PartitionColumnReader::supports_range(range, format));
}
// Formats that can hide physical rows behind deletes are excluded.
for (const auto* table : {"transactional_hive", "iceberg", "paimon"}) {
range.table_format_params.__set_table_format_type(table);
EXPECT_FALSE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_ORC));
}
// Hudi COW carries its partition value in the partition path, exactly like Hive.
range.table_format_params.__set_table_format_type("hudi");
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_PARQUET));
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_ORC));
// A split is fine: the reader no longer needs a row count, so a partial range is no longer a
// correctness problem -- one row per range is the contract.
range.table_format_params.__set_table_format_type("hive");
range.__set_start_offset(1);
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_PARQUET));
range.__set_start_offset(0);
range.__set_size(512);
EXPECT_TRUE(PartitionColumnReader::supports_range(range, TFileFormatType::FORMAT_PARQUET));
}

TEST(TableFormatReaderTest, PartitionValueEmitsOneRowPerRangeAndFillsTypedNulls) {
auto value_slot_desc = create_slot_descriptor(0, "part", TPrimitiveType::INT);
auto null_slot_desc = create_slot_descriptor(1, "null_part", TPrimitiveType::INT, true);
SlotDescriptor value_slot(value_slot_desc);
SlotDescriptor null_slot(null_slot_desc);
std::unordered_map<std::string, uint32_t> block_index {{"part", 0}, {"null_part", 1}};

// The footer row count is irrelevant: a range emits its partition values whether or not the
// file turns out to hold rows. (A partition with no file at all emits nothing, simply because
// it produces no scan range.) Keeping the zero-row case here pins that contract down: an
// implementation that re-gated on "proven nonempty" would fail it.
for (const int64_t footer_rows : {0, 1, 10000}) {
auto inner = std::make_unique<MockTableFormatReader>();
auto* inner_ptr = inner.get();
inner->set_fill_col_name_to_block_idx(&block_index);
inner->set_partition_value("part", "42", &value_slot);
inner->set_partition_value("null_part", "", &null_slot, true);
PartitionColumnReader reader(std::move(inner));
EXPECT_EQ(reader.get_push_down_agg_type(), TPushAggOp::type::PARTITION_VALUE);

Block block;
block.insert(
{value_slot.get_empty_mutable_column(), value_slot.get_data_type_ptr(), "part"});
block.insert(
{null_slot.get_empty_mutable_column(), null_slot.get_data_type_ptr(), "null_part"});
size_t read_rows = 0;
bool eof = false;
ASSERT_TRUE(reader.get_next_block(&block, &read_rows, &eof).ok());
EXPECT_EQ(read_rows, 1) << "footer_rows=" << footer_rows;
EXPECT_EQ(block.rows(), 1);
EXPECT_TRUE(eof);
ASSERT_TRUE(block.check_type_and_column().ok());
EXPECT_EQ(block.get_by_position(0).column->get_int(0), 42);
EXPECT_TRUE(block.get_by_position(1).column->is_null_at(0));

ASSERT_TRUE(reader.get_next_block(&block, &read_rows, &eof).ok());
EXPECT_EQ(read_rows, 0);
EXPECT_EQ(block.rows(), 0);
EXPECT_TRUE(eof);
EXPECT_EQ(inner_ptr->read_calls, 0);
ASSERT_TRUE(reader.close().ok());
EXPECT_EQ(inner_ptr->close_calls, 1);
}
}

} // namespace doris
Loading
Loading