Skip to content

feat(query): add opt-in distributed broadcast join planner - #9404

Draft
discord9 wants to merge 38 commits into
GreptimeTeam:mainfrom
discord9:feat/dist-join-planner
Draft

discord9 wants to merge 38 commits into
GreptimeTeam:mainfrom
discord9:feat/dist-join-planner

Conversation

@discord9

@discord9 discord9 commented Sep 29, 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)

Stacked on #9206; merge #9206 first. This branch includes its review follow-up and read-preference wire correction by ordinary merge.

What's changed and what's your intention?

Add an opt-in, default-off experimental_dist_join path for supported two-scan INNER equijoins. Broadcast the smaller input to the larger table's datanodes rather than joining both remote scans on the frontend.

  • Put SMALL on DataFusion's left/hash-build side and BIG on the local right/probe side; restore SQL column order with an existing Projection. BIG-region routing/pruning and producer identity remain intact. DataFusion's JoinSelection remains in control; there is no forced global input order.
  • Supply ordinary MergeScan row estimates from selected unique regions' leader heartbeat observations, using checked arithmetic. Rows are Inexact; memory and column statistics remain unknown. The disk-size placement heuristic is not a hash-memory estimate. Unknown/unusable statistics leave planning unchanged.
  • Retain complete-result, schema and execution-role tests through the actual SQL setting. No shuffle/general distributed join, self-join or aggregate-side support is added. No new service, protobuf schema or persisted format is introduced. The prerequisite carries typed read preference through supporting receivers using an existing context-extension map.

Limits: Query admission is unchanged. Positive max_concurrent_queries can cause nested requests to time out while outer streams hold permits; default-disabled tests do not establish finite-quota success. Typed read preference is preserved across updated receivers; old receivers still ignore its internal carrier. This change does not add enterprise bootstrap or follower lifecycle support. Earlier release/ECS measurements used an older BIG-build implementation and are not performance evidence for this SMALL-build stack.

Verification: Local checks passed on clean head 2c016aa0ca4fca6ab862b84aabcb0402d0a04863:

  • cargo check --locked -p query -p operator -p session -p meta-client -p datanode -p flow -p servers -p tests-integration --tests — passed.
  • cargo nextest run --locked -p query -p operator -p session -p common-session -p meta-client --lib --retries 0 -j 4 --no-fail-fast — 1,098 passed, 2 existing skips.
  • Focused datanode tests (factory installation and unchanged parallelism) — 2 passed. cargo nextest run --locked -p tests-integration --lib --retries 0 -j 1 --no-fail-fast -E 'test(nested_merge_scan_capability_test) | test(test_flow_requester_health_check_and_handle_request)' — 6 passed, including actual SQL planning/execution and both datanode/Flow strict-ingress regressions.
  • cargo nextest run --locked -p servers --lib --retries 0 -j 2 -E 'test(grpc::greptime_handler::tests) | test(http::hints::tests) | test(http::read_preference::tests)' — 7 passed.
  • cargo clippy --locked -p query -p operator -p session -p meta-client -p datanode -p flow -p servers -p substrait@1.4.0-alpha.0 -p tests-integration --all-targets -- -D warnings — passed. Formatting, TOML, license-header, spelling and Snafu checks passed.
  • Rebuilt greptime and sqlness-runner; greptime --version reports the exact commit above and clean: true. sqlness-runner bare -t '^(standalone|distributed):(create|view|dist_join)$' -c tests/cases --bins-dir "$CARGO_TARGET_DIR/debug" — all 7 distinct standalone/distributed cases passed. Generated SQL expectations are unchanged.

The actual SQL integration retains all 4,877 projected rows against the frontend and known-row references, original SELECT * schema/results, SMALL hash-build/BIG probe roles on each dispatched datanode, and BIG-region pruning. The SQLness case checks off→on→off plans and the complete ordered six-column result. Current-head CI and performance are not yet claimed; Draft-skipped jobs will not be counted as functional test passes.

Additional local enterprise validation: The current #9206/#9404 OSS stack was tested against current enterprise sources with a local companion factory registration. The enterprise-enabled nested/wire/Region/Flow suite passed all 10 selected tests without retries, including normal pre-build-hook construction, real outer/inner TCP Flight follower reads with complete known/frontend result parity, failure without an available follower and no nested leader fallback, and the four supported typed-preference protobuf roundtrips. Strict all-target lint, focused caller/native-datanode regressions, formatting, and a fresh enterprise binary build passed. A separate real enterprise CLI metasrv/datanode/frontend run passed ordinary SQL with the exact complete three-row schema/results; startup identity and owned-process cleanup were checked. The private companion is local only, not part of these OSS PRs or an enterprise release. The focused fixture opens physical followers explicitly; neither it nor the ordinary-SQL CLI check proves automatic follower lifecycle or a single CLI follower end-to-end scenario.

PR Checklist

  • 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).

@github-actions github-actions Bot added docs-required This change requires docs update. size/XXL labels Sep 29, 2026
Give a datanode the ability to decode and execute a plan whose MergeScan
carries a nested plan (e.g. an inner MergeScan over a build table), so a
distributed sub-plan pushed to a datanode can itself fan out to other
datanodes. This is the execution/codec capability only; it is not wired
to any user-facing switch and produces no plan by itself.

Codec (src/query/src/query_engine/default_serializer.rs):
- MergeScanAwareSerializer keeps the request's SessionState and the
  engine catalog manager, and decodes a MergeScan payload's tables by
  (catalog, schema, table) name through the engine catalog rather than
  the region-aware request catalog list (which binds every table to the
  request's region and breaks on mismatched columns).

Datanode (src/datanode):
- New DatanodeRegionQueryHandler (region_query.rs) serves region queries
  for the MergeScan nodes of a received plan, resolving region leaders
  via PartitionRuleManager and dispatching to peer datanodes.
- Wire KvBackendCatalogManager + PartitionRuleManager + NodeClients into
  the datanode's catalog/cache registry (with_dist_planner stays false:
  a datanode never rewrites its own plan).

Bug fixes uncovered while wiring this up:
- cache registry: register the table/table_info/table_name/table_route/
  partition_info caches the datanode's catalog manager and partition
  rule manager require (previously missing -> panic on lookup).
- default_serializer: decode_sub_plan bridge no longer drives a future
  on a current-thread runtime from within another executor (spawn a
  dedicated thread + runtime when not on a multi-thread runtime).

Tested by query_engine::default_serializer unit tests: nested MergeScan
round-trip and engine-catalog resolution with mismatched column names.
No user-facing behavior change.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
… state

SessionStateBuilder::new_from_existing collects the function registry
(HashMap) into a Vec and build() re-registers each entry; register_udf
registers aliases before the function's own name and the Vec order is
unstable. Rebuilding the session state to install MergeScanAwareSerializer
after the GreptimeDB functions were registered could therefore rebind a
name/alias (e.g. date_format, an alias of DataFusion's to_char) back to
the built-in function. This path runs for every decoded plan, not only
nested MergeScan ones.

Extract the GreptimeDB function registration into register_greptime_functions
and re-apply it after the final build of both the top-level decode state and
the MergeScan payload state. Add regression tests asserting date_format stays
bound to the GreptimeDB implementation on both paths.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Three follow-ups from review of the nested MergeScan execution capability:

1. Layer the datanode caches (P2b). table_cache derives from
   table_info_cache + table_name_cache, and partition_info_cache derives
   from table_route_cache, but they were all registered in a single flat
   CacheRegistry whose invalidate() runs concurrently with no ordering.
   A derived cache could be invalidated and refilled from a base cache
   that had not been invalidated yet, leaving stale metadata. Build the
   base caches and the derived caches as separate layers of a
   LayeredCacheRegistry (reusing the existing layering mechanism) so base
   caches finish invalidating before derived ones. Both the production
   path (cmd) and the test scaffolding now build the same layered
   registry; the flat build_datanode_cache_registry is removed.

2. Don't let a nested MergeScan's internal region query compete for the
   concurrency permit (P2a). An outer query holds its permit until its
   stream is fully consumed, and a nested MergeScan makes the datanode
   issue region queries to peer datanodes as execution stages of that
   same query. If such an inner stage also took a permit, the outer query
   waiting on it would deadlock against its own held permit (or time out
   when max_concurrent_queries > 0). Mark inner DN->DN region queries as
   internal stages: QueryRequest gains an  flag carried over
   the wire in the request header's query_context extensions (reserved
   key, not injectable via hints), and the region server skips acquiring
   a concurrency permit for internal-stage requests on both acquire
   sites. Outer queries still acquire normally.

3. Add an end-to-end capability test
   (tests-integration/.../nested_merge_scan_capability_test.rs): hand-build
   Join(<local probe region>, MergeScan(build regions on both datanodes))
   with mismatched column names and duplicate build keys, ship it to a
   datanode over a real TCP Flight endpoint, consume the stream fully and
   compare the row multiset against the frontend's reference join result,
   assert a real cross-datanode DoGet via an RPC counter, and assert an
   inner region failure fails the whole query. Also migrate the
   with_real_datanode_grpc_addr + DatanodeRpcStats test scaffolding.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
…on fork

ScalarUDFImpl's trait bounds changed so its impl can no longer be
downcast via as_any(); assert the function equals the GreptimeDB
implementation instead, which to_char can never satisfy.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Add an end-to-end regression test for the nested MergeScan capability
under max_concurrent_queries = 1: both datanodes' only permit is held by
the outer queries' streams while their inner MergeScans fetch the build
regions from the same datanodes over real Flight connections. The result
consumption is bounded by a timeout, so a broken link anywhere in the
chain (marking the inner request internal, carrying the marker in the
request header, honoring it at the receiving acquire sites) fails the
test instead of deadlocking.

Also add a minimal with_datanode_options_override hook to the cluster
builder, and fix the test module doc: the outer request is handed to
RegionServer::handle_remote_read in process; only the inner region
queries go over the wire.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
The table route cache uses InitStrategy::VersionChecked, whose version
counter belongs to the whole CacheContainer rather than to a single
table: any TableId invalidation bumps the shared version, so a cold
load of one table's route retries when an unrelated table is
invalidated concurrently. Document this as a known limitation on the
region query handler, and pin the behavior with an integration test:
run the nested plans while a background task invalidates an unrelated
table id on every datanode, assert the results still match the expected
multiset (retries only add work, they never return a stale route), and
report the observed load amplification via the cache miss counter
(2 loads without invalidations, 8 over three runs with 24 unrelated
invalidations in this run). Also expose the datanode cache registries
on the test cluster so a test can drive the local invalidation path
that the heartbeat handler drives after a metasrv broadcast.

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>
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>
…join

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
…ss support

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 force-pushed the feat/dist-join-planner branch from e9763f6 to 694489f Compare October 4, 2026 09:50
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>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
…tures

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>
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>
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>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

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

docs-required This change requires docs update. size/XXL

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant