Skip to content

Add soft/hard annotation system - #761

Draft
gabotechs wants to merge 1 commit into
mainfrom
gabrielmusat/soft-hard-task-count
Draft

gabotechs wants to merge 1 commit into
mainfrom
gabrielmusat/soft-hard-task-count

Conversation

@gabotechs

@gabotechs gabotechs commented Oct 2, 2026 •

Copy link
Copy Markdown
Collaborator

Closes #741

Changes the TaskCountAnnotation from:

/// Annotation attached to a single [ExecutionPlan] that determines how many distributed tasks
/// it should run on.
#[derive(Clone, Copy)]
pub enum TaskCountAnnotation {
    /// The desired number of distributed tasks for this node. The final task count for the
    /// annotated node might not be exactly this number, it is more like a hint, so depending
    /// on the desired task count of adjacent nodes, the final task count might change. Fractional
    /// values are preserved while hints are combined and rounded up when a concrete task count is
    /// required.
    Desired(f64),
    /// Sets a maximum number of distributed tasks for this node. Typically used with the inner
    /// value of 1, stating that this node cannot be executed in a distributed fashion.
    Maximum(usize),
}

To

/// Annotation attached to a single [ExecutionPlan] that determines how many distributed tasks
/// it should run on.
#[derive(Clone, Copy)]
pub struct TaskCountAnnotation {
    /// The load of a node measured in tasks required to properly execute it. This number is used
    /// as a hint for the distributed planner to decide on a final task count for a stage.
    /// This value is reconciled with other [TaskCountAnnotation]s provided by other nodes in the
    /// same stage.
    pub soft: f64,
    /// Exact number of tasks that should be allocated to the node.
    pub hard: Option<usize>,
}

So that it can:

  • Represent an annotation that contributes a "load" factor during the task count reconciliation step while still requiring a hard task count.
  • Ensure there's a way of guaranteeing an exact task count for certain nodes, failing at planning time if the expected task count cannot be met.

@gabotechs
gabotechs force-pushed the gabrielmusat/soft-hard-task-count branch from c354452 to 0c34d2e Compare October 2, 2026 11:00
@gabotechs
gabotechs added this pull request to stack #762 October 2, 2026 13:15
@gabotechs
gabotechs force-pushed the gabrielmusat/soft-hard-task-count branch 4 times, most recently from 1e0de92 to a8d25f4 Compare October 4, 2026 08:29
Comment on lines -104 to -106
/// Per-child allocation hint passed to [`ChildrenIsolatorUnionExec::from_children_and_weights`].
#[derive(Debug, Clone, Copy)]
pub(crate) struct ChildWeight {

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We can afford to drop this intermediate struct now and just rely on the existing TaskCountAnnotation

Comment thread tests/work_unit_feed.rs
Comment on lines -484 to +485
┌───── Stage 1 ── tasks=3, partitions=6
│ DistributedUnionExec: t0:[c0(0/2)] t1:[c0(1/2)] t2:[c1]
┌───── Stage 1 ── tasks=2, partitions=6
│ DistributedUnionExec: t0:[c0(0/2), c1] t1:[c0(1/2)]

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is actually a fix:

c1 here refers to this part of the SQL query in the test:

SELECT 'static' as tag, 0 as task, 0 as partition, 'x' as letter

Which is essentially a 1-row in memory table. It's not worth it to allocate a dedicated task just for executing this 1-row in-memory table remotely, it should be completely fine to just execute it in t0 along whatever else happens here, as it's trivial.

@gabotechs
gabotechs force-pushed the gabrielmusat/soft-hard-task-count branch from a8d25f4 to ab5ff8c Compare October 5, 2026 05:42
@gabotechs

Copy link
Copy Markdown
Collaborator Author

benchmarks run tpch/sf100

@gabot-0

gabot-0 commented Oct 5, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 6b937fb3ee95 → PR head ab5ff8c994db · View exact source diff

=== Comparing tpch/sf100 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=1234.0, new=1234.0, diff=no change (sum of per-query averages)
TOTAL: prev=45920 ms, new=47115 ms, diff=1.03 slower ✖
Show full query output
      q1: prev=1644 ms, new=1671 ms, diff=1.02 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1224 ms, new=1237 ms, diff=1.01 slower ✖, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=1967 ms, new=2159 ms, diff=1.10 slower ✖, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 835 ms, new=1013 ms, diff=1.21 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2588 ms, new=2590 ms, diff=1.00 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 947 ms, new= 928 ms, diff=1.02 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2818 ms, new=3029 ms, diff=1.07 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3297 ms, new=3241 ms, diff=1.02 faster ✔, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4129 ms, new=4253 ms, diff=1.03 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4075 ms, new=4146 ms, diff=1.02 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 778 ms, new= 888 ms, diff=1.14 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1298 ms, new=1219 ms, diff=1.06 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 981 ms, new= 929 ms, diff=1.06 faster ✔, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1135 ms, new=1174 ms, diff=1.03 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2546 ms, new=2617 ms, diff=1.03 slower ✖, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 728 ms, new= 661 ms, diff=1.10 faster ✔, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3109 ms, new=3061 ms, diff=1.02 faster ✔, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3361 ms, new=3508 ms, diff=1.04 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1397 ms, new=1425 ms, diff=1.02 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=1850 ms, new=2011 ms, diff=1.09 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4604 ms, new=4724 ms, diff=1.03 slower ✖, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 609 ms, new= 631 ms, diff=1.04 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 223 captured both immutable revisions when the request was queued. The bot fetched and checked out each full commit SHA in detached HEAD, then built and deployed the datafusion-distributed-remote-worker --bin worker target from that checkout.

Identity PR base PR head
Source commit 6b937fb3ee9546c2e3c51ae30b7853a58fa22b45 ab5ff8c994dbe3d8f8f417aba0bc7583de801a8a
Phase Base PR head
Build and deployment 1m 50s 1m 34s
All benchmarks 4m 50s 5m 3s
Benchmark tpch/sf100 4m 50s 5m 3s

Workload: tpch/sf100 · all queries · 1 warmup + 5 measured iterations per query for both revisions

Capacity: 12 c5n.4xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 13m 27s

How to use the benchmark bot

Post a comment whose first non-empty line is:

benchmarks run <suite>/<variant>... [--instance-type <type>] [--nodes <count>] [--iterations <count>] [--base main] [--config <key=value>]...

For example:

benchmarks run tpch/sf10 tpch/sf100 --instance-type m5.2xlarge --nodes 24 --iterations 20 --base main --config distributed.collect_dynamic_filters=false

Currently available datasets: clickbench/0-100, tpcds/sf1, tpch/sf1, tpch/sf10, tpch/sf100. Request one or more, without duplicates. The bot validates availability before provisioning and runs every query in each dataset.

Option Default Supported values and behavior
--instance-type <type> c5n.4xlarge One of the supported instance types listed below, subject to availability within the cluster's availability zones and quota.
--nodes <count> 12 An integer from 1 to 60. The deployment uses one benchmark worker per node.
--iterations <count> 5 A positive safe integer. Measured iterations per query for both revisions; warmup is excluded.
--base main PR base Compare against a snapshot of main; no other explicit base is supported.
--config <key=value> none Apply a safe DataFusion session setting to the PR head only. Repeat for distinct keys; spaces and shell syntax are not supported.

Supported instances and per-worker limits:

Instance type EC2 capacity Worker requests and limits
c5n.2xlarge 8 vCPU, 21 GiB 7 vCPU, 17Gi
c5n.4xlarge 16 vCPU, 42 GiB 15 vCPU, 38Gi
m5.2xlarge 8 vCPU, 32 GiB 7 vCPU, 28Gi
m5.4xlarge 16 vCPU, 64 GiB 15 vCPU, 60Gi
r5.2xlarge 8 vCPU, 64 GiB 7 vCPU, 60Gi
r5.4xlarge 16 vCPU, 128 GiB 15 vCPU, 124Gi

Each node runs one worker. The worker allocation reserves 1 vCPU and 4 GiB for Kubernetes and system processes, then makes the rest of the selected instance available to the benchmark.

Limits: Only authorized users can enqueue jobs. Jobs run serially, the queue holds 20 active jobs, and each requester may have 3. Query selection is not supported. By default, every query uses 1 warmup and 5 measured iterations for both revisions.

@gabotechs

Copy link
Copy Markdown
Collaborator Author

benchmarks run tpch/sf100

@gabot-0

gabot-0 commented Oct 5, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 6b937fb3ee95 → PR head ab5ff8c994db · View exact source diff

=== Comparing tpch/sf100 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=1234.0, new=1234.0, diff=no change (sum of per-query averages)
TOTAL: prev=46891 ms, new=47864 ms, diff=1.02 slower ✖
Show full query output
      q1: prev=1683 ms, new=1716 ms, diff=1.02 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1280 ms, new=1433 ms, diff=1.12 slower ✖, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2034 ms, new=2217 ms, diff=1.09 slower ✖, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 948 ms, new=1088 ms, diff=1.15 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2723 ms, new=2915 ms, diff=1.07 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev=1013 ms, new= 961 ms, diff=1.05 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2863 ms, new=2881 ms, diff=1.01 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3384 ms, new=3245 ms, diff=1.04 faster ✔, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4349 ms, new=4199 ms, diff=1.04 faster ✔, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=3983 ms, new=4156 ms, diff=1.04 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 807 ms, new= 808 ms, diff=1.00 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1254 ms, new=1317 ms, diff=1.05 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 959 ms, new= 939 ms, diff=1.02 faster ✔, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1228 ms, new=1142 ms, diff=1.08 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2651 ms, new=2593 ms, diff=1.02 faster ✔, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 660 ms, new= 641 ms, diff=1.03 faster ✔, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3187 ms, new=3090 ms, diff=1.03 faster ✔, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3437 ms, new=3544 ms, diff=1.03 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1364 ms, new=1423 ms, diff=1.04 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=1849 ms, new=1938 ms, diff=1.05 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4625 ms, new=4889 ms, diff=1.06 slower ✖, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 610 ms, new= 729 ms, diff=1.20 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 224 captured both immutable revisions when the request was queued. The bot fetched and checked out each full commit SHA in detached HEAD, then built and deployed the datafusion-distributed-remote-worker --bin worker target from that checkout.

Identity PR base PR head
Source commit 6b937fb3ee9546c2e3c51ae30b7853a58fa22b45 ab5ff8c994dbe3d8f8f417aba0bc7583de801a8a
Phase Base PR head
Build and deployment 1m 43s 1m 35s
All benchmarks 4m 57s 5m 5s
Benchmark tpch/sf100 4m 57s 5m 5s

Workload: tpch/sf100 · all queries · 1 warmup + 5 measured iterations per query for both revisions

Capacity: 12 c5n.4xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 13m 28s

How to use the benchmark bot

Post a comment whose first non-empty line is:

benchmarks run <suite>/<variant>... [--instance-type <type>] [--nodes <count>] [--iterations <count>] [--base main] [--config <key=value>]...

For example:

benchmarks run tpch/sf10 tpch/sf100 --instance-type m5.2xlarge --nodes 24 --iterations 20 --base main --config distributed.collect_dynamic_filters=false

Currently available datasets: clickbench/0-100, tpcds/sf1, tpch/sf1, tpch/sf10, tpch/sf100. Request one or more, without duplicates. The bot validates availability before provisioning and runs every query in each dataset.

Option Default Supported values and behavior
--instance-type <type> c5n.4xlarge One of the supported instance types listed below, subject to availability within the cluster's availability zones and quota.
--nodes <count> 12 An integer from 1 to 60. The deployment uses one benchmark worker per node.
--iterations <count> 5 A positive safe integer. Measured iterations per query for both revisions; warmup is excluded.
--base main PR base Compare against a snapshot of main; no other explicit base is supported.
--config <key=value> none Apply a safe DataFusion session setting to the PR head only. Repeat for distinct keys; spaces and shell syntax are not supported.

Supported instances and per-worker limits:

Instance type EC2 capacity Worker requests and limits
c5n.2xlarge 8 vCPU, 21 GiB 7 vCPU, 17Gi
c5n.4xlarge 16 vCPU, 42 GiB 15 vCPU, 38Gi
m5.2xlarge 8 vCPU, 32 GiB 7 vCPU, 28Gi
m5.4xlarge 16 vCPU, 64 GiB 15 vCPU, 60Gi
r5.2xlarge 8 vCPU, 64 GiB 7 vCPU, 60Gi
r5.4xlarge 16 vCPU, 128 GiB 15 vCPU, 124Gi

Each node runs one worker. The worker allocation reserves 1 vCPU and 4 GiB for Kubernetes and system processes, then makes the rest of the selected instance available to the benchmark.

Limits: Only authorized users can enqueue jobs. Jobs run serially, the queue holds 20 active jobs, and each requester may have 3. Query selection is not supported. By default, every query uses 1 warmup and 5 measured iterations for both revisions.

@gabotechs

Copy link
Copy Markdown
Collaborator Author

benchmarks run tpch/sf100

@gabot-0

gabot-0 commented Oct 5, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 6b937fb3ee95 → PR head ab5ff8c994db · View exact source diff

=== Comparing tpch/sf100 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=1234.0, new=1234.0, diff=no change (sum of per-query averages)
TOTAL: prev=46366 ms, new=46886 ms, diff=1.01 slower ✖
Show full query output
      q1: prev=1656 ms, new=1773 ms, diff=1.07 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1253 ms, new=1242 ms, diff=1.01 faster ✔, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2002 ms, new=1990 ms, diff=1.01 faster ✔, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev=1041 ms, new= 950 ms, diff=1.10 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2658 ms, new=2634 ms, diff=1.01 faster ✔, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 981 ms, new=1018 ms, diff=1.04 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2820 ms, new=2945 ms, diff=1.04 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3215 ms, new=3273 ms, diff=1.02 slower ✖, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4050 ms, new=4147 ms, diff=1.02 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4012 ms, new=4020 ms, diff=1.00 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 819 ms, new=1025 ms, diff=1.25 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1263 ms, new=1310 ms, diff=1.04 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 927 ms, new= 935 ms, diff=1.01 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1181 ms, new=1161 ms, diff=1.02 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2605 ms, new=2617 ms, diff=1.00 slower ✖, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 659 ms, new= 666 ms, diff=1.01 slower ✖, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3141 ms, new=3149 ms, diff=1.00 slower ✖, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3406 ms, new=3521 ms, diff=1.03 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1491 ms, new=1405 ms, diff=1.06 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=1902 ms, new=1908 ms, diff=1.00 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4640 ms, new=4551 ms, diff=1.02 faster ✔, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 644 ms, new= 646 ms, diff=1.00 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 225 captured both immutable revisions when the request was queued. The bot fetched and checked out each full commit SHA in detached HEAD, then built and deployed the datafusion-distributed-remote-worker --bin worker target from that checkout.

Identity PR base PR head
Source commit 6b937fb3ee9546c2e3c51ae30b7853a58fa22b45 ab5ff8c994dbe3d8f8f417aba0bc7583de801a8a
Phase Base PR head
Build and deployment 2m 29s 1m 34s
All benchmarks 4m 54s 5m 8s
Benchmark tpch/sf100 4m 54s 5m 8s

Workload: tpch/sf100 · all queries · 1 warmup + 5 measured iterations per query for both revisions

Capacity: 12 c5n.4xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 14m 13s

How to use the benchmark bot

Post a comment whose first non-empty line is:

benchmarks run <suite>/<variant>... [--instance-type <type>] [--nodes <count>] [--iterations <count>] [--base main] [--config <key=value>]...

For example:

benchmarks run tpch/sf10 tpch/sf100 --instance-type m5.2xlarge --nodes 24 --iterations 20 --base main --config distributed.collect_dynamic_filters=false

Currently available datasets: clickbench/0-100, tpcds/sf1, tpch/sf1, tpch/sf10, tpch/sf100. Request one or more, without duplicates. The bot validates availability before provisioning and runs every query in each dataset.

Option Default Supported values and behavior
--instance-type <type> c5n.4xlarge One of the supported instance types listed below, subject to availability within the cluster's availability zones and quota.
--nodes <count> 12 An integer from 1 to 60. The deployment uses one benchmark worker per node.
--iterations <count> 5 A positive safe integer. Measured iterations per query for both revisions; warmup is excluded.
--base main PR base Compare against a snapshot of main; no other explicit base is supported.
--config <key=value> none Apply a safe DataFusion session setting to the PR head only. Repeat for distinct keys; spaces and shell syntax are not supported.

Supported instances and per-worker limits:

Instance type EC2 capacity Worker requests and limits
c5n.2xlarge 8 vCPU, 21 GiB 7 vCPU, 17Gi
c5n.4xlarge 16 vCPU, 42 GiB 15 vCPU, 38Gi
m5.2xlarge 8 vCPU, 32 GiB 7 vCPU, 28Gi
m5.4xlarge 16 vCPU, 64 GiB 15 vCPU, 60Gi
r5.2xlarge 8 vCPU, 64 GiB 7 vCPU, 60Gi
r5.4xlarge 16 vCPU, 128 GiB 15 vCPU, 124Gi

Each node runs one worker. The worker allocation reserves 1 vCPU and 4 GiB for Kubernetes and system processes, then makes the rest of the selected instance available to the benchmark.

Limits: Only authorized users can enqueue jobs. Jobs run serially, the queue holds 20 active jobs, and each requester may have 3. Query selection is not supported. By default, every query uses 1 warmup and 5 measured iterations for both revisions.

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

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Introduce Exact task count hints

2 participants