feature(stats) Estimate scan cardinality from selected partitions for OLAP and external tables - #64032
feature(stats) Estimate scan cardinality from selected partitions for OLAP and external tables#64032foxtail463 wants to merge 2 commits into
Conversation
|
Thank you for your contribution to Apache Doris. Please clearly describe your PR:
|
|
/review |
ad131d6 to
6b3af18
Compare
There was a problem hiding this comment.
Automated review summary for PR 64032
I reviewed the partition-selection refactor across Nereids logical/physical OLAP and file scans, partition pruning, predicate post-processing, statistics, MV/HBO helpers, and Hive/Hudi/MaxCompute scan translation/runtime paths. I did not find a concrete blocking correctness issue to raise as an inline comment.
Critical checkpoint conclusions:
- Goal and tests: The PR centralizes partition-selection state and extends selected-partition row-count handling for external scans. The changed tests cover logical/physical propagation, pruning behavior, and stats cases relevant to the refactor.
- Scope/minimality: The change is broad but focused on replacing scattered selected-partition/prune flags with
PartitionSelectionabstractions and updating required call sites. - Concurrency/lifecycle: No new shared mutable state, lock ordering, or lifecycle-sensitive static initialization issue was found. New selection objects are immutable snapshots.
- Config/compatibility: No new configuration, storage format, edit log, or FE-BE protocol compatibility change was found.
- Parallel paths: OLAP, Hive/Hudi file scans, MaxCompute, physical translation, stats, and query partition collection were updated consistently.
- Conditional checks: New conditions around partition-pruned/constraint state match prior semantics after tracing the old fields and runtime scan-node consumers.
- Test coverage: Unit tests were added/updated for the main behavior touched. I did not run the test suite in this review environment.
- Observability: Existing logging/profile paths remain sufficient for the changed planning/statistics paths; no missing mandatory metric/log was identified.
- Transaction/persistence/data writes: The PR is planning/statistics-side only; no transaction or persisted metadata path change was identified.
- Performance: The immutable snapshot approach avoids repeated mutable state hazards; no obvious pathological new hot-path scan or allocation issue was confirmed.
User focus points: No additional user-provided review focus was present.
|
run buildall |
6b3af18 to
3a2808c
Compare
|
run buildall |
TPC-H: Total hot run time: 29223 ms |
TPC-DS: Total hot run time: 168995 ms |
FE Regression Coverage ReportIncrement line coverage |
3a2808c to
a4484e0
Compare
|
run buildall |
FE UT Coverage ReportIncrement line coverage |
TPC-H: Total hot run time: 28444 ms |
TPC-DS: Total hot run time: 168570 ms |
FE Regression Coverage ReportIncrement line coverage |
|
/review |
There was a problem hiding this comment.
Automated review summary for PR #64032:
I reviewed the full FE/Nereids diff and the existing inline thread context. I did not find a new blocking correctness issue to raise in this pass.
Critical checkpoint conclusions:
- Goal fit: The partition-selection refactor consistently carries selected partitions into OLAP/file scans and uses selected-partition row counts for stats where the row count is known.
- Correctness: I checked the main semantic paths: OLAP partition pruning, file/Hive/Hudi/MaxCompute partition selection, physical translation, block-rule signaling, and filter selectivity skipping. The retained physical partition predicate for DELETE/skip-prune remains semantically redundant after partition pruning, so I did not treat that as a defect.
- Plan/state propagation: Copy methods and logical-to-physical translators generally preserve
PartitionSelection; places that intentionally narrow selected partitions keep the applied-predicate snapshot only for subsets. - Compatibility and persistence: No storage format, transaction log, FE-BE protocol, or compatibility issue found in the reviewed changes.
- Concurrency/lifecycle: No new locking, async lifecycle, or resource-management risk found.
- Security: This PR does not appear to change authentication, authorization, tenancy, or network-boundary behavior.
- Tests: The PR adds focused unit coverage for selected-partition stats and predicate selectivity skipping. I did not run tests in this review runner.
User focus: .code-review.Uvi60a/review_focus.txt did not contain additional focus points, so there were no extra focus-specific findings.
|
PR approved by anyone and no changes requested. |
a4484e0 to
08f1e9c
Compare
|
run buildall |
TPC-H: Total hot run time: 29556 ms |
TPC-DS: Total hot run time: 174057 ms |
ClickBench: Total hot run time: 25.4 s |
FE UT Coverage ReportIncrement line coverage |
FE Regression Coverage ReportIncrement line coverage |
|
run nonConcurrent |
FE Regression Coverage ReportIncrement line coverage |
|
run nonConcurrent |
FE Regression Coverage ReportIncrement line coverage |
morrySnow
left a comment
There was a problem hiding this comment.
Review Summary
This PR introduces a PartitionSelection abstraction (65 files, +1424/-1075) to fix double-counting of partition predicates in Nereids row-count estimation. The core idea — compute scan row count from selected partitions, then tag Statistics with already-applied conjuncts so computeFilter can skip them — is sound and well-executed. The refactoring is comprehensive: all old-API callers are updated, all withXxx copy methods correctly carry the new state, and the physical plan translator wires all scan node types correctly.
Below are the issues found, ranked by severity:
1. [HIGH] MV/optStats early-return misses conjunctsAppliedToRowCount
File: fe/fe-core/src/main/java/org/apache/doris/nereids/stats/StatsCalculator.java (around line 560)
The optStats.isPresent() early-return path returns derivedStats (from optStats.get()) directly without setting conjunctsAppliedToRowCount. The two other exit paths (partition-pruned and table-level) both set this field via the builder. When a Filter sits above this scan, computeFilter will re-apply the partition predicate selectivity to the already-pruned row count, producing a double underestimate.
Concrete failure: MV scan with partition pruning and stored optStats. The scan row count already reflects partition selection (e.g., 10M of 100M total). The Filter above applies the partition predicate again (×0.1 selectivity), reducing to 1M. The true row count should remain ~10M. This underestimation can cause wrong join order and aggregation strategy.
Suggested fix: Before returning, wrap in a StatisticsBuilder:
return new StatisticsBuilder(derivedStats)
.setConjunctsAppliedToRowCount(
olapScan.getPartitionSelection()
.getAppliedPartitionConjuncts(olapScan, ((Plan) olapScan).getOutput()))
.build();2. [MEDIUM] buildSlotReplaceMap throws Preconditions.checkState — old code returned null gracefully
File: fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/algebra/PartitionSelection.java (new file)
The old PrunePartitionPredicate.buildSlotReplaceMap (deleted) returned null when a snapshot slot had no matching output slot, and the caller checked if (slotReplaceMap != null) to skip the entry. The new PartitionSelection.buildSlotReplaceMap throws IllegalStateException, which would crash query planning if column pruning removes a partition column from the scan output (the column is only referenced in WHERE, not SELECT).
Suggested fix: Restore the null-safe behavior.
3. [MEDIUM] computeFilter removes conjuncts by expression identity — fragile after rewrites
File: fe/fe-core/src/main/java/org/apache/doris/nereids/stats/StatsCalculator.java (computeFilter method)
conjuncts.removeAll(inputStats.getConjunctsAppliedToRowCount()) uses Expression.equals() to match. If expression canonicalization produces a syntactically different but semantically equivalent form, removeAll silently fails to match, causing double-counting with no warning.
4. [LOW] LinkedHashSet allocation on every computeFilter call
File: fe/fe-core/src/main/java/org/apache/doris/nereids/stats/StatsCalculator.java (computeFilter method)
Even when conjunctsAppliedToRowCount is empty (the common case), every call allocates a new LinkedHashSet. Guarding with isEmpty() would avoid overhead for the dominant path.
5. [LOW] PartitionPruner no longer records always-TRUE partition predicates
File: fe/fe-core/src/main/java/org/apache/doris/nereids/rules/expression/rules/PartitionPruner.java
When a partition predicate folds to BooleanLiteral.TRUE, the old code recorded it so the post-processor could strip it from the plan. The new code returns Optional.empty(), so these trivially-true predicates survive to BE execution — a minor optimization regression.
6. [LOW] isPruned() method vs partitionPruned field creates two-check inconsistency
Files: ExternalPartitionSelection.java, OlapPartitionSelection.java
Some callers check partitionSelection.partitionPruned directly, others call partitionSelection.isPruned(). These disagree when partitionPruned=true but all partitions survive. For OlapPartitionSelection, there is no isPruned() method at all — the equivalent check is done ad-hoc in computeOlapScan.
|
run buildall |
4a8cd6a to
2f5157a
Compare
|
run buildall |
TPC-H: Total hot run time: 27802 ms |
TPC-DS: Total hot run time: 152045 ms |
ClickBench: Total hot run time: 23.92 s |
FE UT Coverage ReportIncrement line coverage |
… OLAP and external tables
2f5157a to
38cb694
Compare
|
run buildall |
TPC-H: Total hot run time: 27971 ms |
TPC-DS: Total hot run time: 152341 ms |
ClickBench: Total hot run time: 23.81 s |
FE UT Coverage ReportIncrement line coverage |
FE Regression Coverage ReportIncrement line coverage |
Problem Summary
After partition pruning the optimizer still derived scan cardinality from table-level row counts, so external-table scans kept the unpruned table cardinality and OLAP scans double-applied partition predicates that remained in the plan for materialized-view rewrite, skewing join orderand cost estimates.
Solution
Add table-level selected-partition row count APIs (OLAP distributes unreported partition rows over the remaining partitions; external tables estimate from file listings within the selected set) and consume them in StatsCalculator, which now records the prunable partition conjuncts as conjunctsAppliedToRowCount so computeFilter no longer double-counts them while still narrowing their column domains against table-level statistics; the recorded proof is carried across scan slot rebinding and finally consumed by the PrunePartitionPredicate post-processor after MV rewrite, and external column statistics are proportionally scaled to the narrowed scan cardinality