Skip to content

Stream and partition GpuAggregateExec instead of collecting the whole input (audit P4) #29

Description

@vyncint

Today — GpuAggregateExec collects its entire input into one batch before aggregating:

// crates/oxidelake-compute/src/exec/aggregate.rs:209-214
let stream = futures::stream::once(async move {
    let batches = collect(input, context).await?;
    let all = concat_batches(&input_schema, &batches)?;
    ...
    let out = operator.aggregate(&all, &spec)?;

with Partitioning::UnknownPartitioning(1) (:84); the planner inserts CoalescePartitionsExec under it for multi-partition sources (crates/oxidelake-planner/src/rule.rs:366-370). The CUDA path adds u32_len(input.num_rows, "aggregate input")? and a group table of 2 × rows slots (crates/oxidelake-device/src/cuda/ops.rs:772-773). The DataFusion plan it replaces was a streaming, partitioned, two-phase Partial/Final aggregate. docs/audit-2026-09-02.md:55 records this as deferred P4.

Why it is worth fixing — every GROUP BY on a cuda/metal-target session, and every aggregate on a cluster with OXIDE_CLUSTER_BACKEND=cuda, materialises the whole input on one node and runs single-threaded; inputs above 4,294,967,295 rows error. GPU placement is a scalability regression for the most common analytical query. A production milestone cannot ship GPU-targeted aggregation in this shape.

Fix — per-batch typed accumulators keyed by group, merged at the end, keeping DataFusion's Partial/Final split so Ballista parallelises the partial half (lower only the Partial half and leave the stock Final, or lower both with a partitioned exchange between them). Until then, gate the rewrite behind a row-count estimate so small inputs keep the fast path and large ones stay on DataFusion.

Done when —

  • cluster_sql_matches_embedded_sql runs its 11 MiB multi-partition fixture with the aggregate lowered and EXPLAIN shows more than one partition feeding the GPU aggregate.
  • Peak memory of a 10M-row GROUP BY k under --target cuda (embedded, CPU fallback) is bounded by batch size, measured with /usr/bin/time -v and recorded.
  • The u32::MAX rows limit applies per batch, not per query.

Activity

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

    bugSomething isn't working

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions