Visitar URL original
Adapative Query Runner (AQE) PoC · Issue #19485 · apache/pinot · GitHub
Skip to content

Adapative Query Runner (AQE) PoC #19485

Description

@wirybeaver

Why AQE? Let runtime data change the remaining plan

SET adaptiveExecution=true;
SET numGroupsLimit=100000;
SELECT col1, COUNT(*) FROM a GROUP BY col1;

Suppose the plan assigns eight workers to the final GROUP BY, but the upstream aggregation produces only a few KiB across eight HASH partitions. Running all eight consumers wastes scheduling and operator state. AQE-PR-5 can use the observed size to coalesce those partitions into one consumer, retaining every original partition without rehashing it.

initial plan:   producer group -> 8 HASH partitions -> 8 consumers -> broker
runtime:       producers finish -> stats/manifests -> AQE rules / replan
adaptive plan:  same 8 partitions -----------------> 1 consumer  -> broker

Coalescing is the first rule, not the framework's limit. Ordered rules can call a physical replanner and return a replacement pending plan, changing operators, stages, routing and parallelism—not just a worker count.

Another example: skip work when the build is empty

SELECT /*+ joinOptions(join_strategy='dynamic_broadcast') */ o.order_id
FROM orders o
WHERE o.customer_id IN (
  SELECT c.customer_id FROM customers c WHERE c.region = 'new-region'
);

Suppose orders is large, but the customer filter produces zero rows at runtime. The result must be empty: preparing and executing probe-side SSE requests is wasted work. AQE-PR-1 uses that runtime fact to skip those requests, while preserving the normal MSE EOS and statistics path.

Query Execution Flow

The two branches below are complementary, independent execution paths.

                         SQL -> MSE plan
                                |
                     run producer/build work
                                |
             +------------------+----------------------+
             |                                         |
   dynamic-broadcast build                  selected HASH exchange
     (existing breaker)                                |
             |                           files + handles + row/byte stats
   empty build? [AQE-PR-1]                         [AQE-PR-2]
        /          \                                   |
      yes           no                  all workers succeed + streams close
       |             |                            [AQE-PR-3]
   skip probe     normal probe                          |
   SSE requests   SSE execution            snapshot -> ordered AQE rules
        \          /                              [AQE-PR-5]
         result/EOS                           /                 \
                                    coalesce M -> K     other rule/replanner*
                                              \                 /
                                         validate + publish pending plan
                                                  [AQE-PR-5]
                                                       |
                                            rebuild ready stage groups
                                             [AQE-PR-3 + AQE-PR-5]
                                                       |
                                          bind original partition handles
                                             [AQE-PR-4 + AQE-PR-5]
                                                       |
                                           dispatch ready consumer group
                                                  [AQE-PR-3]
                                                       |
                                           read assigned partition files
                                             [AQE-PR-2 + AQE-PR-4]
                                                       |
                                                broker reduction

* Extension point: PR5 ships coalescing, not skew-splitting or join rules. Replanning changes only unsubmitted work; completed stages and the output contract stay fixed. Each new dispatch group must fit the original thread budget, while an individual stage may grow.

Implementation Roadmap

PR Deliverable Review
AQE-PR-1 Skip probe SSE work after an empty dynamic-broadcast build. apache/pinot#19470
AQE-PR-2 Materialized HASH partitions with runtime row/byte statistics. wirybeaver/pinot#70
AQE-PR-3 Completion barriers and dependency-ordered dispatch. wirybeaver/pinot#71
AQE-PR-4 Late-bound partition handles for undispatched consumers. wirybeaver/pinot#72
AQE-PR-5 Rule-driven runtime replanning, with deterministic partition coalescing (M -> K). wirybeaver/pinot#73

PR1 is independent. Review the adaptive path as an incremental fork stack:

aqe/base -> AQE-PR-2 -> AQE-PR-3 -> AQE-PR-4 -> AQE-PR-5

Upstream PR2–4 remain temporarily closed in favor of these incremental reviews. PR2–4 alone keep worker counts fixed; PR5 adds opt-in adaptation through adaptiveExecution=true. The coalescing target defaults to 64 MiB (aqeTargetPartitionBytes); 0 disables only that rule. The PoC coalesces eligible unary HASH consumers while respecting group caps and local limits. Skew splitting and join rules remain follow-ups—not framework restrictions.

Activity

added
queryRelated to query processing
PEP-RequestPinot Enhancement Proposal request to be reviewed.
multi-stageRelated to the multi-stage query engine
on Sep 17, 2026

wirybeaver commented on Oct 4, 2026

@wirybeaver
ContributorAuthor

Full-query end-to-end evidence: AQE-PR-4 and AQE-PR-5

These use a two-server fixture with real broker/server gRPC, query operators and materialized partition reads. Results below are local fixture evidence, not production-cluster benchmarks.

AQE-PR-4: 2 cases passed

Test target Observed result
testMaterializedStagedQuery: nonempty GROUP BY Fixed parallelism; result rows/schema match the nonmaterialized baseline; input handles are late-bound.
testMaterializedStagedQuery: empty GROUP BY The empty result/schema matches baseline through the complete staged path.

These integrate PR2 materialized exchange + PR3 staged dispatch + PR4 late binding. PR4 test instructions and source.

AQE-PR-5: 3 new cases + the 2 PR4 regressions = 5 cases passed

New test target Observed result
testAdaptiveCoalescingExecutesAllOriginalPartitions: nonempty GROUP BY M -> 1; every original partition handle accounted for; results/schema unchanged; no missing or merge-failed worker reports.
The same method: empty GROUP BY One consumer; the same empty result/schema; complete worker coverage.
testRuleReplansEmptyInputIntoBrokerValues A second, test-only rule generates a replacement broker plan that actually executes; the removed consumer is not dispatched; unused materialized files are cleaned up.

The two inherited cases are exactly testMaterializedStagedQuery with nonempty and empty output above. The PR5 run at 54cd91b914 passed all 5, with 0 failures; it is not five newly added PR5 cases. PR5 test instructions, five-case breakdown and source.

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    PEP-RequestPinot Enhancement Proposal request to be reviewed.multi-stageRelated to the multi-stage query enginequeryRelated to query processing

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions