Repository navigation
feat(datanode): execute nested MergeScan plans on the datanode - #9206
Merged
killme2008 merged 21 commits intoOct 10, 2026
Merged
killme2008 merged 21 commits into
killme2008 merged 21 commits into
Conversation
discord9
force-pushed
the
feat/datanode-nested-mergescan-capability
branch
2 times, most recently
from
September 18, 2026 04:25
f26bd58 to
9e0addc
Compare
discord9
marked this pull request as ready for review
September 18, 2026 06:40
Contributor
discord9
force-pushed
the
feat/datanode-nested-mergescan-capability
branch
from
September 18, 2026 09:25
9e0addc to
25d8dd7
Compare
discord9
force-pushed
the
feat/datanode-nested-mergescan-capability
branch
from
September 28, 2026 03:13
c63c599 to
08bb9f1
Compare
5 of 6 tasks
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
force-pushed
the
feat/datanode-nested-mergescan-capability
branch
from
October 4, 2026 09:13
08bb9f1 to
2a46309
Compare
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
killme2008
reviewed
Oct 8, 2026
killme2008
reviewed
Oct 8, 2026
killme2008
left a comment
Member
There was a problem hiding this comment.
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.
…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>
killme2008
reviewed
Oct 8, 2026
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
approved these changes
Oct 10, 2026
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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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
MergeScannodes. This PR supplies the capability; it does not enable distributed rewriting or generate nested plans for ordinary queries.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.RegionQueryHandlerFactoryRef, with the OSS leader handler as fallback. Streaming clients retain an unset request timeout.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_queriescan 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, andgit 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 head6fcfd44b5fbf7b4b40e181cfc981b6cdd212b6e0: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, andgit diff --check— passed.Previous head
6fcfd44b5fbCI 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
backport-<target>labels (e.g.backport-v1.3targetsrelease/v1.3).