Skip to content

feat(datanode): execute nested MergeScan plans on the datanode - #9206

Merged
killme2008 merged 21 commits into
GreptimeTeam:mainfrom
discord9:feat/datanode-nested-mergescan-capability
Oct 10, 2026
Merged

killme2008 merged 21 commits into
GreptimeTeam:mainfrom
discord9:feat/datanode-nested-mergescan-capability

Conversation

@discord9

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

Prerequisite for #9404, which separately chooses broadcast-join plans.

What's changed and what's your intention?

Allow datanodes to decode and execute received plans containing nested MergeScan nodes. This PR supplies the capability; it does not enable distributed rewriting or generate nested plans for ordinary queries.

  • Decode nested payloads asynchronously using the registered request SessionState. Outer scans stay region-bound; payload scans use the engine catalog with request defaults and context. Lookup errors do not fall back to the request table.
  • Wire the datanode's catalog, partition manager and region clients through the existing RegionQueryHandlerFactoryRef, with the OSS leader handler as fallback. Streaming clients retain an unset request timeout.
  • Supply the catalog/routing caches this path needs, using existing base-before-derived invalidation layers. The protobuf schema and placeholder semantics are unchanged.
  • Preserve the typed read preference through received query contexts using one reserved internal extension. Missing legacy metadata stays Leader; malformed or unsupported present values are rejected before execution. Default Leader does not add a carrier, and client-supplied generic hints cannot overwrite it.

The integration tests hand-build nested plans, compare complete duplicate-sensitive results with frontend SQL and known rows, verify cross-datanode Flight requests, and require missing/failed inner regions to fail the query. Review follow-up removes obsolete runtime-bridge guards, redundant cold-route/cache-presence tests and their exclusive harness, the unused DataFusion information-schema override, and misleading comments. One missing-table negative regression remains; the redundant denied-access case and its exclusive test scaffold have been removed.

Limits: Admission is unchanged. A positive max_concurrent_queries can cause outer streams to hold permits needed by nested requests; acquisition may time out and fail the query. Default-disabled tests do not establish success under finite quotas. Enterprise bootstrap must install its follower-aware factory. Read-preference preservation requires receivers with this context decoder; old receivers still ignore the new extension. The carrier does not add enterprise bootstrap or follower lifecycle support. No performance or concurrent cache-invalidation guarantee is claimed.

Verification: Local checks passed on clean tests-only follow-up cbb7bf68f74baba91e3559ea5633b1c472b167be, which removes the redundant negative case without changing production code:

  • cargo check --locked -p query --tests — passed.
  • cargo nextest run --locked -p query --lib --retries 0 -j 6 --no-fail-fast — 807 passed, 2 existing skips; the retained missing-table regression passed.
  • cargo clippy --locked -p query --all-targets -- -D warnings, cargo fmt --all -- --check, and git diff --check — passed.

Current-head CI completed successfully on cbb7bf68f74: Rust CI, Integration CI, Checks, and dependency checks. The current-head rollup is 55 successful checks and 2 skipped checks (coverage and recent-release compatibility), with no pending or failed checks. The following previously completed local checks were run on clean head 6fcfd44b5fbf7b4b40e181cfc981b6cdd212b6e0:

  • cargo check --locked -p tests-integration --tests — 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)' — 4 passed, including both datanode ingress paths and Flow invalid-context rejection.
  • cargo nextest run --locked -p session -p common-session --lib --retries 0 -j 2 — 16 passed, including the frozen legacy-shaped protobuf fixture.
  • 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 session -p datanode -p flow -p servers -p tests-integration --all-targets -- -D warnings, cargo fmt --all -- --check, taplo format --check src/session/Cargo.toml, and git diff --check — passed.

Previous head 6fcfd44b5fb CI completed successfully: Rust CI, Integration CI, Checks, Docs and dependency checks. That head’s PR rollup had 54 successful checks, 2 skipped checks (coverage and recent-release compatibility), and one cancelled semantic-title check. The frozen compatibility fixture verifies the new decoder's handling of legacy metadata; it does not show that old receivers honor the new carrier.

Additional local enterprise validation: The #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 companion factory registration has now been submitted for enterprise review as a Draft PR; it is not merged or released. Its merge order is #9206, the normal upstream-main submodule update, then companion revalidation. These combined-source checks are separate from the companion PR checkout, which still uses the unchanged earlier OSS pin. 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 size/M docs-not-required This change does not impact docs. size/XXL and removed size/M labels Sep 17, 2026
@discord9
discord9 force-pushed the feat/datanode-nested-mergescan-capability branch 2 times, most recently from f26bd58 to 9e0addc Compare September 18, 2026 04:25
@discord9
discord9 marked this pull request as ready for review September 18, 2026 06:40
@discord9
discord9 requested review from a team, evenyag and v0y4g3r as code owners September 18, 2026 06:40
@discord9
discord9 force-pushed the feat/datanode-nested-mergescan-capability branch from 9e0addc to 25d8dd7 Compare September 18, 2026 09:25
@discord9
discord9 force-pushed the feat/datanode-nested-mergescan-capability branch from c63c599 to 08bb9f1 Compare September 28, 2026 03:13
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>
@discord9
discord9 force-pushed the feat/datanode-nested-mergescan-capability branch from 08bb9f1 to 2a46309 Compare October 4, 2026 09:13
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Comment thread src/datanode/src/region_query.rs
Comment thread src/query/src/query_engine/default_serializer.rs Outdated

@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.

The capability path looks correct. Most of the remaining diff is tests and harness code that guard against designs no longer in this PR; trimming it would roughly halve the review surface.

Before #9404 lets anyone turn this on, please open an issue for nested-request admission under a positive max_concurrent_queries.

The description could be much shorter. The cache-registry "bug" can't happen on main, which uses DummyCatalogManager and never looks those caches up. The date_format issue came from an earlier revision of this PR. The verification logs belong in CI.

Comment thread tests-integration/src/tests/nested_merge_scan_capability_test.rs Outdated
Comment thread src/query/src/query_engine/default_serializer.rs Outdated
Comment thread src/query/src/query_engine/default_serializer.rs Outdated
Comment thread src/datanode/src/region_server.rs Outdated
Comment thread src/datanode/src/datanode.rs Outdated
Comment thread src/cache/src/lib.rs Outdated
Comment thread src/cache/src/lib.rs Outdated
Comment thread tests-integration/src/standalone.rs Outdated
Comment thread src/datanode/src/region_query.rs Outdated
Comment thread tests-integration/src/tests/nested_merge_scan_capability_test.rs Outdated
…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>
Comment thread src/query/src/query_engine/default_serializer.rs Outdated
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>
@killme2008
killme2008 added this pull request to the merge queue Oct 10, 2026
Merged via the queue into GreptimeTeam:main with commit 28697c9 Oct 10, 2026
59 checks passed
discord9 added a commit to discord9/greptimedb that referenced this pull request Oct 10, 2026
Adapt the production and test serializer callers added by GreptimeTeam#9206 and remove the obsolete partition_cols assertion.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

docs-not-required This change does not impact docs. size/XXL

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants