Repository navigation
Conversation
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>
|
Warning @discord9 has 7 open non-draft pull requests in this repository, over the limit of 5. Review is the scarcest resource here. Please land or close some of these before
This check is advisory for now and blocks nothing. |
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+--------+---------------------+" |
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
Question: why we chose hash distribution before?
killme2008
left a comment
There was a problem hiding this comment.
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) % 4routing 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'vsk < '4'(expected4 | 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?
| 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"#; |
There was a problem hiding this comment.
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() { |
There was a problem hiding this comment.
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(), |
There was a problem hiding this comment.
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.
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, includingINsubqueries with remote dynamic-filter pushdown.MergeScanExecdistributes regions across output partitions by their position in the selected region vector. It does not hash rows by the storage partition columns. However, it advertisedPartitioning::Hash, andPassDistributioncould 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:
UnknownPartitioningwith the existing output partition count.PassDistributionrule. Existing DataFusion distribution enforcement inserts real hash repartitioning when needed.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.IN/left-semi regression with parallelism 2 and remote dynamic filters explicitly enabled, preserving all four keys and their earliest timestamps.UnionExecwith per-input sorts rather than claiming interleave-compatible hash buckets. The complete UNION ALL result still contains1, 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, 20instead of5,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=Partitioneddirectly over differently orderedMergeScanExecinputs, without actual repartitioning. At this host's default parallelism 20, real hash repartitioning was inserted and all 12 runs returned5000 | 5000; this masked the defect rather than fixing the distribution contract. Parallelism 1 also returned correct results withCollectLeft.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, andgit 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 realHash(k)repartitioning on both JOIN inputs even though the region orders differ; parallelism 1 usesCollectLeft.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 keys1, 11, 21, 31with earliest timestamp1791504000000in 10/10 runs. The parallelism-1, RDF-off, and conditional-aggregate controls also returned all four expected rows.EXPLAIN ANALYZE VERBOSEconfirms a partitioned left-semi join with realHash(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 TQLexplainin 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 twoInterleaveExec→UnionExecnames and three obsoletephysical_plan after PassDistributionRulerows. 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.
backport-<target>labels (e.g.backport-v1.3targetsrelease/v1.3).