Skip to content

fix(query): stop reporting region scans as hash partitioned - #9491

Open
discord9 wants to merge 6 commits into
GreptimeTeam:mainfrom
discord9:fix/9456-mergescan-partitioning
Open

discord9 wants to merge 6 commits into
GreptimeTeam:mainfrom
discord9:fix/9456-mergescan-partitioning

Conversation

@discord9

@discord9 discord9 commented Oct 9, 2026 •

Copy link
Copy Markdown
Contributor

I hereby agree to the terms of the GreptimeDB CLA.

Refer to a related PR or issue link (optional)

Fixes #9456.
Fixes #9492.
Related to #8460, which handled joins using only a subset of storage partition columns.

What's changed and what's your intention?

Fix silent missing matches in partitioned hash joins over MergeScanExec, including IN subqueries with remote dynamic-filter pushdown.

MergeScanExec distributes regions across output partitions by their position in the selected region vector. It does not hash rows by the storage partition columns. However, it advertised Partitioning::Hash, and PassDistribution could change the advertised hash expressions without redistributing any rows. A partitioned hash join then compared corresponding output partition indices even when those indices represented different regions or key ranges.

This change:

  • Reports UnknownPartitioning with the existing output partition count.
  • Removes the metadata-only distribution rewrite and the PassDistribution rule. Existing DataFusion distribution enforcement inserts real hash repartitioning when needed.
  • Keeps logical partition aliases and aggregate pushdown, scan parallelism, ordering inference, and region execution unchanged.
  • Adds an executed physical LEFT JOIN regression with deliberately permuted peers, three regions striped over two output partitions, an unmatched left row, and preservation of remote dynamic-filter producer IDs beneath the real shuffles.
  • Adds standalone and distributed integration coverage for the reported root LEFT JOIN (child GROUP BY k) shape: duplicate append-mode children, complete ordered results, a genuinely unmatched root, parallelism 16/8, and asymmetric region pruning.
  • Adds an exact eight-row, four-region IN/left-semi regression with parallelism 2 and remote dynamic filters explicitly enabled, preserving all four keys and their earliest timestamps.
  • Refreshes the runner-generated UNION plan snapshot: unknown scan partitioning uses UnionExec with per-input sorts rather than claiming interleave-compatible hash buckets. The complete UNION ALL result still contains 1, 2, 3, 3.

Reproduction and trigger

Using the issue's exact 16-region table and seed-1 fixture (5,000 roots and 15,000 children), an isolated pre-fix standalone build returned incorrect results in all 24 runs at query parallelism 16/8. An independent six-run verification at parallelism 16 observed matching counts 0, 33, 16, 0, 0, 20 instead of 5,000; all 5,000 roots remained present. Source counts and independent aggregate/UNION controls confirmed the stored data was intact.

The bad plans used HashJoinExec: mode=Partitioned directly over differently ordered MergeScanExec inputs, without actual repartitioning. At this host's default parallelism 20, real hash repartitioning was inserted and all 12 runs returned 5000 | 5000; this masked the defect rather than fixing the distribution contract. Parallelism 1 also returned correct results with CollectLeft.

Issue #9492 exercises another consequence of the same distribution mismatch. Its two scans have the same region order, but a partitioned semi join creates a dynamic filter routed by hash(key) % N. Region-position stripes are not those hash buckets, so the filter can reject matching rows even when the scan order is aligned. The correction preserves remote dynamic filtering and makes both join inputs genuinely hash partitioned rather than using RDF-off as the fix.

The fix may add shuffles to plans that previously relied on the incorrect hash claim. This is necessary to preserve complete query results; it does not alter persisted data, SQL semantics, or public configuration.

Validation

Production/query verification and actual product runs use the implementation at 6a6dd0b578aa963f0fbf616775237907f3a226dd; later follow-ups add the semi-join integration regression and refresh runner-generated EXPLAIN snapshots without changing production Rust:

  • cargo fmt --all -- --check, hawkeye check, python3 scripts/check-snafu.py, and git diff --check.

  • cargo check --locked -p query --tests.

  • cargo nextest run --locked -p query --lib --no-fail-fast: 798 passed, 2 existing skips. This includes the executed physical JOIN regression.

  • cargo clippy --locked -p query --all-targets -- -D warnings.

  • cargo check --locked -p tests-integration --tests.

  • cargo nextest run --locked -p tests-integration --lib -E 'test(test_partitioned_merge_scan_)' -j 1 --no-fail-fast: 4 passed, covering both LEFT JOIN and RDF-enabled semi join in standalone and distributed. These instance tests belong to the library harness.

  • cargo build --locked -p cmd --bin greptime -p sqlness-runner --bin sqlness-runner. The fresh binary reports the clean implementation commit above.

  • Original issue, actual HTTP SQL: recreated the exact seed-1 5,000-key fixture, including identical child rows sharing each root's timestamp. At parallelism 16/8/1/20, six runs each all returned 5000 | 5000 (24/24). For every setting, the full ordered result was compared against all 5,000 expected keys with child count 3, with no LIMIT or response truncation. Source counts remained 5,000 roots / 15,000 children. Raw plans at 16/8 confirm real Hash(k) repartitioning on both JOIN inputs even though the region orders differ; parallelism 1 uses CollectLeft.

  • Focused SQLness, actual standalone/distributed product: ran the prebuilt runner with --case-dir <isolated selected-case copies> --bins-dir <fresh build> --test-filter '^(standalone|distributed):(pass_distribution_partition_subset_join|subqueries|step_aggr_advance|join_filter_pushdown_edge|tsid_binary_join_regression|partition|tsid_column)$' --jobs 1 --preserve-state. 9/9 passed on the verification run, with byte-stable regenerated results. All non-EXPLAIN outputs remained unchanged. Eight snapshots were identical; the one UNION EXPLAIN change was reviewed and copied from the runner output, not hand-edited. This includes the fix: repartition subset partition key joins #8460 subset-key join's complete four cross-matches and the PromQL binary-join cases in both modes.

  • Issue IN subquery returns incomplete results with remote dynamic filter pushdown #9492, actual HTTP SQL: recreated its exact eight-row, four-range-partition table. Default parallelism 2 and explicit query_parallelism=2, both with RDF enabled, each returned all keys 1, 11, 21, 31 with earliest timestamp 1791504000000 in 10/10 runs. The parallelism-1, RDF-off, and conditional-aggregate controls also returned all four expected rows. EXPLAIN ANALYZE VERBOSE confirms a partitioned left-semi join with real Hash(key_id, 2) repartitioning on both inputs, the same peer order on both scans, and active dynamic filters; the RDF-off/parallelism-1 controls change the expected plan paths.

  • Additional CI-reported SQLness plans: directly regenerated and reviewed tql-cte, order_by_exceptions, and shared TQL explain in standalone/distributed. 6/6 passed on a second execution with byte-stable expectations. All returned data and error text remained byte-identical; the only updates are two InterleaveExec→UnionExec names and three obsolete physical_plan after PassDistributionRule rows. Both modes generate identical shared results.

The PR remains Ready for review. For head f10e9987d37c1ed988a6c96f7703d328419263f9, Rust CI, Checks, Integration CI, and Check Dependencies all completed successfully. This includes Rust tests, Clippy, Windows tests/build, riscv64 build, all five SQLness configurations, and the integration/fuzz/chaos jobs. Coverage and recent-release compatibility jobs were skipped, not counted as passes. The preceding Integration CI failures are not counted as passes; their three additional shared plan expectations were regenerated and verified before this successful current-head run.

PR Checklist

Please convert it to a draft if some of the following conditions are not met.

  • I have written the necessary rustdoc comments.
  • I have added the necessary unit tests and integration tests.
  • This PR requires documentation updates.
  • API changes are backward compatible.
  • Schema or data changes are backward compatible.
  • This PR needs to be backported to release branches, and I have added the backport-<target> labels (e.g. backport-v1.3 targets release/v1.3).

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
@discord9
discord9 marked this pull request as ready for review October 9, 2026 10:39
@discord9
discord9 requested review from a team and evenyag as code owners October 9, 2026 10:39
@github-actions

github-actions Bot commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
.await;
assert_eq!(
output,
"+--------+---------------------+\n| key_id | first_ts |\n+--------+---------------------+\n| 1 | 2026-10-09T00:00:00 |\n| 11 | 2026-10-09T00:00:00 |\n| 21 | 2026-10-09T00:00:00 |\n| 31 | 2026-10-09T00:00:00 |\n+--------+---------------------+"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Prefer:

let table = r#"
+--------+---------------------+
| key_id | first_ts            |
+--------+---------------------+
| 1      | 2026-10-09T00:00:00 |
| 11     | 2026-10-09T00:00:00 |
| 21     | 2026-10-09T00:00:00 |
| 31     | 2026-10-09T00:00:00 |
+--------+---------------------+
"#;

.collect();
let partitioning = Partitioning::Hash(partition_exprs, output_partition_count);
// Region stripes are positional, not hash buckets; do not advertise a hash distribution.
let partitioning = Partitioning::UnknownPartitioning(output_partition_count);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question: why we chose hash distribution before?

@killme2008 killme2008 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please add the two issue reproductions as a sqlness case under tests/cases/standalone/common/, with an EXPLAIN that pins mode=Partitioned and RepartitionExec on both inputs. A later change that flips the join to CollectLeft would keep the result-only tests green while no longer exercising this path.

sqlness runs at parallelism = 4 with RDF on and has no per-query hint today:

  • #9492: add a few more keys per range so hash(key) % 4 routing reliably disagrees with the region stripes.
  • #9456: one key per region plus an unmatched root, and a pruned pair with equal region counts but shifted stripes, e.g. k >= '1' AND k < '5' vs k < '4' (expected 4 | 3 | 9).

Please confirm both fail on a pre-fix build. A -- SQLNESS HINT query_parallelism=2 interceptor in the runner would let #9492 run at its original parallelism; that can be a separate PR.

question: not blocking, but this puts a hash shuffle on both sides of every PromQL binary op between metric-engine tables. #7927 relied on the old claim (map_hash_requirement_through_projection) to skip it, and no sqlness case covers a partitioned PromQL join, so the cost is unmeasured. Can you post one before/after number for rate(a) / rate(b) on a partitioned physical table, and open a follow-up issue for a real partition-wise join when both scans share the same region list?

Comment on lines +382 to +383
FROM (SELECT k FROM join_repro WHERE kind = 'root' AND k < '8') r
LEFT JOIN (SELECT k, count(*) n FROM join_repro WHERE kind = 'child' AND k < 'c' GROUP BY k) c ON c.k = r.k"#;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These bounds keep the stripes aligned: at parallelism 8 region i lands in partition i on both sides, and at 16 the partition counts differ (8 vs 12), so the old code shuffled too. This case likely passes without the fix. Shift one side by a region so both sides have 8 regions but misaligned stripes (exercised at parallelism 8):

FROM (SELECT k FROM join_repro WHERE kind = 'root' AND k >= '1' AND k < '9') r
LEFT JOIN (SELECT k, count(*) n FROM join_repro WHERE kind = 'child' AND k < '8' GROUP BY k) c ON c.k = r.k"#;

Expected becomes 8 | 7 | 21.

}

#[tokio::test]
async fn merge_scan_left_join_shuffles_unknown_partitioning_before_hash_join() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This passes on the old code too: the helper passes an empty AliasMapping, so the scan advertised Hash([], 2), which EnsureRequirements already treats as unsatisfied. The UnknownPartitioning assertion near L1574 is what guards the fix. Drop the execution and producer-id parts, or build the scan so the old code would have skipped the shuffle.

self.region_query_handler.clone(),
query_ctx,
session.config().target_partitions(),
merge_scan.partition_cols().clone(),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: this was the last reader of MergeScanLogicalPlan::partition_cols (merge_scan.rs:341, getter at :464). Drop the field and the getter, or leave a follow-up.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

2 participants