Skip to content

WIP: tune - #743

Draft
Rich-T-kid wants to merge 1 commit into
datafusion-contrib:mainfrom
Rich-T-kid:rich-T-kid/tune-config-mn-product
Draft

Rich-T-kid wants to merge 1 commit into
datafusion-contrib:mainfrom
Rich-T-kid:rich-T-kid/tune-config-mn-product

Conversation

@Rich-T-kid

Copy link
Copy Markdown
Contributor

see #717 (comment)

base command

benchmarks run tpch/sf1 --config distributed.two_step_shuffle_fanout_threshold=10

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=64

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=256

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base db6a84347ad3PR head 7d5b7e91e5f8 · 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=47311 ms, new=52347 ms, diff=1.11 slower ❌
Show full query output
      q1: prev=1657 ms, new=1787 ms, diff=1.08 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1362 ms, new=1201 ms, diff=1.13 faster ✔, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2195 ms, new=1905 ms, diff=1.15 faster ✔, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 955 ms, new= 940 ms, diff=1.02 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2700 ms, new=2606 ms, diff=1.04 faster ✔, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 914 ms, new= 964 ms, diff=1.05 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=3006 ms, new=2756 ms, diff=1.09 faster ✔, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3216 ms, new=3204 ms, diff=1.00 faster ✔, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4230 ms, new=4144 ms, diff=1.02 faster ✔, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4426 ms, new=4074 ms, diff=1.09 faster ✔, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 792 ms, new= 884 ms, diff=1.12 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1276 ms, new=1310 ms, diff=1.03 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 914 ms, new= 932 ms, diff=1.02 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1157 ms, new=1358 ms, diff=1.17 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2539 ms, new=8265 ms, diff=3.26 slower ❌, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 655 ms, new= 839 ms, diff=1.28 slower ✖, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3053 ms, new=3221 ms, diff=1.06 slower ✖, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3459 ms, new=3387 ms, diff=1.02 faster ✔, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1380 ms, new=1327 ms, diff=1.04 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=2064 ms, new=1962 ms, diff=1.05 faster ✔, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4725 ms, new=4646 ms, diff=1.02 faster ✔, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 636 ms, new= 635 ms, diff=1.00 faster ✔, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 185 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 db6a84347ad3ce8dbb04038d545fa008d9998d9a 7d5b7e91e5f88af6f6fcbb65c807f37850769739
Phase Base PR head
Build and deployment 1m 49s 1m 34s
All benchmarks 5m 0s 5m 34s
Benchmark tpch/sf100 5m 0s 5m 34s

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

Capacity: 12 c5n.4xlarge nodes for both revisions

PR-head configs: distributed.two_step_shuffle_fanout_threshold=64

Other timings: Queue 2s · Dataset validation 0s · Total 14m 8s

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.

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base db6a84347ad3PR head 7d5b7e91e5f8 · 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=48434 ms, new=52591 ms, diff=1.09 slower ✖
Show full query output
      q1: prev=1653 ms, new=1702 ms, diff=1.03 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1276 ms, new=1651 ms, diff=1.29 slower ✖, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2179 ms, new=2142 ms, diff=1.02 faster ✔, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 980 ms, new=1087 ms, diff=1.11 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2702 ms, new=2947 ms, diff=1.09 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 976 ms, new=1001 ms, diff=1.03 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=3182 ms, new=3213 ms, diff=1.01 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3377 ms, new=3614 ms, diff=1.07 slower ✖, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4075 ms, new=4643 ms, diff=1.14 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4332 ms, new=4925 ms, diff=1.14 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 969 ms, new=1015 ms, diff=1.05 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1324 ms, new=1502 ms, diff=1.13 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 947 ms, new=1211 ms, diff=1.28 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1279 ms, new=1455 ms, diff=1.14 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2938 ms, new=2838 ms, diff=1.04 faster ✔, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 667 ms, new= 883 ms, diff=1.32 slower ✖, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3070 ms, new=3318 ms, diff=1.08 slower ✖, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3400 ms, new=3681 ms, diff=1.08 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1402 ms, new=1506 ms, diff=1.07 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=2212 ms, new=2299 ms, diff=1.04 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4823 ms, new=5226 ms, diff=1.08 slower ✖, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 671 ms, new= 732 ms, diff=1.09 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 186 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 db6a84347ad3ce8dbb04038d545fa008d9998d9a 7d5b7e91e5f88af6f6fcbb65c807f37850769739
Phase Base PR head
Build and deployment 1m 36s 1m 34s
All benchmarks 5m 5s 5m 54s
Benchmark tpch/sf100 5m 5s 5m 54s

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

Capacity: 12 c5n.4xlarge nodes for both revisions

PR-head configs: distributed.two_step_shuffle_fanout_threshold=256

Other timings: Queue 14m 11s · Dataset validation 0s · Total 14m 16s

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.

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=96

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base db6a84347ad3PR head 7d5b7e91e5f8 · 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=48243 ms, new=51689 ms, diff=1.07 slower ✖
Show full query output
      q1: prev=1761 ms, new=1674 ms, diff=1.05 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1365 ms, new=1344 ms, diff=1.02 faster ✔, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2342 ms, new=2872 ms, diff=1.23 slower ✖, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 897 ms, new=1351 ms, diff=1.51 slower ❌, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2717 ms, new=4548 ms, diff=1.67 slower ❌, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 933 ms, new=1023 ms, diff=1.10 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2954 ms, new=3099 ms, diff=1.05 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3365 ms, new=4049 ms, diff=1.20 slower ✖, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4295 ms, new=4315 ms, diff=1.00 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4176 ms, new=4195 ms, diff=1.00 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev=1022 ms, new=1113 ms, diff=1.09 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1352 ms, new=1213 ms, diff=1.11 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 939 ms, new= 943 ms, diff=1.00 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1178 ms, new=1332 ms, diff=1.13 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2726 ms, new=2620 ms, diff=1.04 faster ✔, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 703 ms, new= 692 ms, diff=1.02 faster ✔, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3197 ms, new=3262 ms, diff=1.02 slower ✖, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3560 ms, new=3510 ms, diff=1.01 faster ✔, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1378 ms, new=1286 ms, diff=1.07 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=2027 ms, new=1923 ms, diff=1.05 faster ✔, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4719 ms, new=4581 ms, diff=1.03 faster ✔, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 637 ms, new= 744 ms, diff=1.17 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 187 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 db6a84347ad3ce8dbb04038d545fa008d9998d9a 7d5b7e91e5f88af6f6fcbb65c807f37850769739
Phase Base PR head
Build and deployment 1m 38s 1m 35s
All benchmarks 5m 7s 5m 31s
Benchmark tpch/sf100 5m 7s 5m 31s

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

Capacity: 12 c5n.4xlarge nodes for both revisions

PR-head configs: distributed.two_step_shuffle_fanout_threshold=96

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

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.

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=192

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=160

@gabotechs

Copy link
Copy Markdown
Collaborator

@Rich-T-kid reminder that if you don't specify --nodes it will use 12 machines with 15 CPUs each (180 partitions). When submitting commands, you might also want to keep this into account.

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base db6a84347ad3PR head 7d5b7e91e5f8 · 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=47870 ms, new=51555 ms, diff=1.08 slower ✖
Show full query output
      q1: prev=1672 ms, new=1701 ms, diff=1.02 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1458 ms, new=1669 ms, diff=1.14 slower ✖, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2120 ms, new=2179 ms, diff=1.03 slower ✖, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 879 ms, new=1023 ms, diff=1.16 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2585 ms, new=2744 ms, diff=1.06 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 942 ms, new= 897 ms, diff=1.05 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2914 ms, new=3248 ms, diff=1.11 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3439 ms, new=3986 ms, diff=1.16 slower ✖, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4332 ms, new=4668 ms, diff=1.08 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4316 ms, new=4413 ms, diff=1.02 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 816 ms, new=1028 ms, diff=1.26 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1345 ms, new=1425 ms, diff=1.06 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev=1008 ms, new=1144 ms, diff=1.13 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1241 ms, new=1297 ms, diff=1.05 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2812 ms, new=2871 ms, diff=1.02 slower ✖, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 717 ms, new= 869 ms, diff=1.21 slower ✖, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3109 ms, new=3311 ms, diff=1.06 slower ✖, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3467 ms, new=3692 ms, diff=1.06 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1383 ms, new=1453 ms, diff=1.05 slower ✖, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=1955 ms, new=2170 ms, diff=1.11 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4732 ms, new=5046 ms, diff=1.07 slower ✖, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 628 ms, new= 721 ms, diff=1.15 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 188 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 db6a84347ad3ce8dbb04038d545fa008d9998d9a 7d5b7e91e5f88af6f6fcbb65c807f37850769739
Phase Base PR head
Build and deployment 1m 35s 1m 35s
All benchmarks 5m 5s 5m 32s
Benchmark tpch/sf100 5m 5s 5m 32s

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

Capacity: 12 c5n.4xlarge nodes for both revisions

PR-head configs: distributed.two_step_shuffle_fanout_threshold=192

Other timings: Queue 2s · Dataset validation 0s · Total 13m 55s

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.

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base db6a84347ad3PR head 7d5b7e91e5f8 · 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=47409 ms, new=49061 ms, diff=1.03 slower ✖
Show full query output
      q1: prev=1699 ms, new=1625 ms, diff=1.05 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
      q2: prev=1395 ms, new=1228 ms, diff=1.14 faster ✔, tasks: prev=92.0, new=92.0, diff=no change
      q3: prev=2004 ms, new=1985 ms, diff=1.01 faster ✔, tasks: prev=52.0, new=52.0, diff=no change
      q4: prev= 945 ms, new=1042 ms, diff=1.10 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
      q5: prev=2662 ms, new=2663 ms, diff=1.00 slower ✖, tasks: prev=79.0, new=79.0, diff=no change
      q6: prev= 921 ms, new= 914 ms, diff=1.01 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=2944 ms, new=2901 ms, diff=1.01 faster ✔, tasks: prev=79.0, new=79.0, diff=no change
      q8: prev=3366 ms, new=3606 ms, diff=1.07 slower ✖, tasks: prev=94.0, new=94.0, diff=no change
      q9: prev=4273 ms, new=4294 ms, diff=1.00 slower ✖, tasks: prev=100.0, new=100.0, diff=no change
     q10: prev=4516 ms, new=4117 ms, diff=1.10 faster ✔, tasks: prev=65.0, new=65.0, diff=no change
     q11: prev= 790 ms, new= 795 ms, diff=1.01 slower ✖, tasks: prev=65.0, new=65.0, diff=no change
     q12: prev=1241 ms, new=1526 ms, diff=1.23 slower ✖, tasks: prev=48.0, new=48.0, diff=no change
     q13: prev= 906 ms, new=1089 ms, diff=1.20 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q14: prev=1164 ms, new=2184 ms, diff=1.88 slower ❌, tasks: prev=26.0, new=26.0, diff=no change
     q15: prev=2559 ms, new=3387 ms, diff=1.32 slower ✖, tasks: prev=50.0, new=50.0, diff=no change
     q16: prev= 668 ms, new= 658 ms, diff=1.02 faster ✔, tasks: prev=51.0, new=51.0, diff=no change
     q17: prev=3132 ms, new=3087 ms, diff=1.01 faster ✔, tasks: prev=38.0, new=38.0, diff=no change
     q18: prev=3437 ms, new=3447 ms, diff=1.00 slower ✖, tasks: prev=64.0, new=64.0, diff=no change
     q19: prev=1410 ms, new=1264 ms, diff=1.12 faster ✔, tasks: prev=26.0, new=26.0, diff=no change
     q20: prev=1902 ms, new=2033 ms, diff=1.07 slower ✖, tasks: prev=63.0, new=63.0, diff=no change
     q21: prev=4798 ms, new=4556 ms, diff=1.05 faster ✔, tasks: prev=86.0, new=86.0, diff=no change
     q22: prev= 677 ms, new= 660 ms, diff=1.03 faster ✔, tasks: prev=32.0, new=32.0, diff=no change
Verification and run details

Job 189 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 db6a84347ad3ce8dbb04038d545fa008d9998d9a 7d5b7e91e5f88af6f6fcbb65c807f37850769739
Phase Base PR head
Build and deployment 1m 4s 1m 35s
All benchmarks 5m 11s 5m 8s
Benchmark tpch/sf100 5m 11s 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

PR-head configs: distributed.two_step_shuffle_fanout_threshold=160

Other timings: Queue 13m 57s · Dataset validation 0s · Total 13m 4s

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.

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=64 --nodes 60

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark job 195 failed for tpch/sf100. Full details are available in the controller journal.

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.

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=96 --nodes 60

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

benchmarks run tpch/sf100 --config distributed.two_step_shuffle_fanout_threshold=160 --nodes 60

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark job 196 failed for tpch/sf100. Full details are available in the controller journal.

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.

@gabot-0

gabot-0 commented Sep 23, 2026

Copy link
Copy Markdown

Requested by this comment.

Benchmark job 197 failed for tpch/sf100. Full details are available in the controller journal.

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.

3 participants