Skip to content

coordinator: display consumer dynamic filters after execution - #623

Merged
jayshrivastava merged 27 commits into
mainfrom
js/1-display-dynamic-filters
Sep 8, 2026
Merged

jayshrivastava merged 27 commits into
mainfrom
js/1-display-dynamic-filters

Conversation

@jayshrivastava

@jayshrivastava jayshrivastava commented Aug 11, 2026 •

Copy link
Copy Markdown
Collaborator

Stack

This stack of PRs implements distributed dynamic filtering #528

  1. coordinator: display consumer dynamic filters after execution #623 <- you are here
  2. feat: plan distributed dynamic filters #634
  3. feat: forward remote dynamic filter updates to coordinator #635
  4. coordinator: merge partial dynamic filters  #636
  5. coordinator: forward merged dynamic filters to consumers #637
  6. worker: apply merged dynamic filters during execution #639

Closes: #529

Problem

Post df-55 upgrade, dynamic filters should work in the worker-local case. There's no way to observe them working other than looking at metrics.

  ┌───── Stage 2 ── tasks=1
  │ AggregateExec: Final COUNT(*)
  │   [Stage 1] => NetworkCoalesceExec
  └──────────────────────────────────────────────────
    ┌───── Stage 1 ── tasks=2
    │ HashJoinExec: orders.customer_id = selected_customers.customer_id
    │   DistributedLeafExec:
    |     ...
    │   DistributedLeafExec:
    │     t0: DataSourceExec: predicate=DynamicFilter [ empty ]
    │     t1: DataSourceExec: predicate=DynamicFilter [ empty ]
    └────────────────────────────────────────────────

Ideally we want the final filters visible when displaying plans.

Solution

This PR adds a new protocol which is basically identical to the metrics protocol. Even the MetricsStore is now just Store and is generic over TaskMetrics and TaskCompletedDynamicFilters (contains completed dynamic filters for a task).

pub(crate) type MetricsStore = Store<TaskMetrics>;
pub(crate) type CompletedDynamicFilterStore = Store<TaskCompletedDynamicFilters>;

Similar to the metrics protocol, workers now collect completed dynamic filters and send them back to the coordinator.

Coordinator                                               Worker
-----------                                               ------
       Create independent display copies
                    |
                    +-- SetPlan(task 0, filter IDs) -------> Decode plan
                    |                                       |
                    |                                       | execute
                    |                                       |
                    |                                       |
                    |                                       |
                    |                                       |
                    |                                       | task finishes
                    |                                       v
                    |<----- TaskDynamicFilters ----- Serialize completed filters from the consumers
                    |
                    v

Then, at display time, we call apply_reports_to_distributed_leaves which traverses the plan_for_viz and updates the dynamic filters for all the variants:

DistributedLeafExec
  task 0: DynamicFilter [ key@0 >= 1 AND key@0 <= 10 ]
  task 1: DynamicFilter [ empty ]

Notes

Duplicate RPC Messages

We will eventually have more dynamic filter RPCs which manage the worker -> coordinator -> merge -> worker flow mentioned in #553.

In theory, the coordinator will know at merge time what the completed filters are, making the TaskCompletedDynamicFilters and final worker -> coordinator message in this PR irrelevant.

However, I think having these mechanisms be separate is good because a) it helps us validate that the dynamic filter coordinator -> worker flow work using external "oracle", and b) there's no guarantee that the coordinator -> worker propagation happens before the query is done (ex. the DataSourceExec may not block execution waiting for dynamic filters), so it's good to have a separate way to know if the final DataSourceExec applied a filter or not.

AND true and empty filters

DynamicFilter [ sr_returned_date_sk@0 >= 2451545 AND sr_returned_date_sk@0 <= 2451910 AND true ] AND DynamicFilter [ empty ]

In this filter AND true occurs because of apache/datafusion#24277. The first DynamicFilter is active but we lose the HashTableLookupExpr when serializing it to send back to the coordinator.

The 2nd filter is DynamicFilter [ empty ] because this is a dynamic filter produced by a remote producer, which does not get propagated to this node yet. This will be fixed later.

Displaying Dynamic Filters

Protocol is as similar to the metrics protocol as possible. Due to double wrapping (MetricsWrapperExec wraps DistributedLeafExec, it's tricky to do the dynamic filter rewrite after doing the metrics rewrite. So rewrite_distributed_plan_with_dynamic_filters has to be called first.

let plan = rewrite_distributed_plan_with_dynamic_filters(plan).await?;
let plan = rewrite_distributed_plan_with_metrics(plan, DistributedMetricsFormat::Aggregated).await?;
println!("{}", display_plan_ascii(plan.as_ref(), true));

Testing

  • Tests in tests/dynamic_filtering.rs

@jayshrivastava jayshrivastava changed the title display dynamic filters during execution display dynamic filters after execution Aug 11, 2026
@jayshrivastava
jayshrivastava force-pushed the js/1-display-dynamic-filters branch from 2a2bffc to f549dc2 Compare August 13, 2026 13:16
@jayshrivastava
jayshrivastava changed the base branch from js/upgrade-df-55-08-10 to branch-55 August 13, 2026 13:16
@jayshrivastava jayshrivastava changed the title display dynamic filters after execution coordinator: display dynamic filters after execution Aug 13, 2026
@stuhood

stuhood commented Aug 15, 2026

Copy link
Copy Markdown
Contributor

Thanks for working on this: this will be very useful!

One quick thought: the dynamic filter will sometimes be much, much larger than what you would actually want to display in an EXPLAIN plan (a large InList or Hash). That suggests that rather than sending the whole filter, what would actually make sense to send back is some sort of human readable summary of the filter?

Also, I have a draft of a related change on our codebase, and it seemed like the easiest mechanism for transferring this kind of information back is via metrics... but the most natural/obvious thing that seemed to be missing in that case was essentially a "string" metric type (we would use it to display a chosen strategy/enum from a scan). Do you think that that might be worth pursuing upstream?

@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

One quick thought: the dynamic filter will sometimes be much, much larger than what you would actually want to display in an EXPLAIN plan (a large InList or Hash). That suggests that rather than sending the whole filter, what would actually make sense to send back is some sort of human readable summary of the filter?

Also, I have a draft of a related change on our codebase, and it seemed like the easiest mechanism for transferring this kind of information back is via metrics... but the most natural/obvious thing that seemed to be missing in that case was essentially a "string" metric type (we would use it to display a chosen strategy/enum from a scan). Do you think that that might be worth pursuing upstream?

Serializing them as a string is reasonable. Rather than a metric, I think we can implement a PhysicalExpr which just wraps a string and inject it into the play for display using DynamicFilterPhysicalExpr::update(string_expr). @gabotechs what do you think?

@gabotechs

Copy link
Copy Markdown
Collaborator

🤔 I'm not sure if I'm understanding the suggestion. Updating a filter with DynamicFilterPhysicalExpr::update is not really related to visualization, it's how you actually update the filter no?

@jayshrivastava
jayshrivastava force-pushed the js/1-display-dynamic-filters branch 2 times, most recently from 4ccec74 to bac65ba Compare August 18, 2026 18:52
Base automatically changed from branch-55 to main August 20, 2026 08:38
gabotechs added a commit that referenced this pull request Aug 20, 2026
## Summary 

Closes
#530

- This PR updates the upstream datafusion SHA to the HEAD of
https://github.com/apache/datafusion/commits/branch-55/ (edit: this
branch is continuously being updated. I will make sure this PR is at the
head before merging)
- Rust upgade to 1.94


## Changes

1. In `src/protobuf/distributed_codec.rs` we now use the
`proto_converter` argument during serde
- We still don't use the `DeduplicatingProtoConverter`, so dynamic
filters don't necessarily work. I think this is outside the scope of
this PR will be addressed in
#623,
which will be rebased after the upgrade.

3. `ExecutionPlan::apply_expressions` is added for every custom
`ExecutionPlan` in this repo
- Wrapper types (`MetricsWrapperExec`, `WorkUnitFileScanConfig`,
`DistributedLeafExec`) delegate to the inner type
- Other plans take`TreeNodeRecursion::Continue` because they have no
expressions (ex. `SamplerExec`)
- Note that `apply_expressions` does not need to yield sort or
partitioning expressions in the plan properties

3. We migrate from `partition_statistics` to `statistics_from_inputs`
for every `ExecutionPlan`.
- `src/distributed_planner/statistics/plan_statistics.rs` can just use
`statistics_from_inputs` directly instead of doing the
`StatisticsWrapper` workaround.

5. Range partitioning is now supported.
- CPU costing now includes range-key comparison cost and has a new unit
test. See src/distributed_planner/
   statistics/complexity_cpu.rs:238.
- I think there's open questions about range partitioning. I've opened
an issue here to make sure it behaves as expected after the upgrade:
#628 (comment)

6. Peak-memory metrics use the existing gauge wire representation.

DataFusion added MetricValue::PeakMemoryUsage. It is serialized as the
existing named-gauge protobuf variant to avoid a wire-format change. See
src/protocol/grpc/
   metrics_proto.rs:124.

The value and name survive, and aggregation is still additive, but
decoding produces a generic Gauge, not PeakMemoryUsage. The practical
difference is mainly display formatting: it
may render as a count rather than human-readable bytes. This is the
clearest remaining compromise/risk in the upgrade.

7. File-scan rebalancing changed its discriminator.

DataFusion removed partitioned_by_file_group;
output_partitioning.is_some() is now the source of truth. See
src/events/defaults/file_scan_config.rs:43. This decides whether files
are
round-robin rebalanced or split through FileGroupPartitioner, so it is
behavior-sensitive even though it is a one-line migration.

8. Two previously ignored correctness tests were enabled.
- See `tests/multi_task_collect_join_repros.rs`
- These were upstream DataFusion correctness fixes, not fixes made
locally in this upgrade.

9. drop(reporter) was made explicit on the sampler’s empty-input path.

The reporter sends its result on Drop; explicitly dropping it both
satisfies the new compiler/lint behavior and guarantees the zero-row EOS
report is sent before returning. See src/
   execution_plans/sampler.rs:259.

10. Plan changes

- `dynamic_rg_pruning=eligible` is now displayed on eligible scans:
1,354 occurrences in TPC-DS, 188 in TPC-H, and 12 in ClickBench
- `DataSourceExec` now displays its output partitioning. See
`tests/join.rs` (eventually, someone should delete this test
#628)
- Project after sort. This looks like some upstream optimizer rule
change ex. `tests/distributed_unions.rs` and
`tests/distributed_aggregation.rs`.
```
-          │   SortExec: expr=[MinTemp@0 ASC NULLS LAST, RainToday@1 ASC NULLS LAST], preserve_partitioning=[true]
-          │     ProjectionExec: expr=[MaxTemp@0 as MinTemp, RainToday@1 as RainToday]
+          │   ProjectionExec: expr=[MaxTemp@0 as MinTemp, RainToday@1 as RainToday]
+          │     SortExec: expr=[MaxTemp@0 ASC NULLS LAST, RainToday@1 ASC NULLS LAST], preserve_partitioning=[true]
```
- LocalLimitExec became more common: TPC-DS went from 0 to 20
occurrences and ClickBench from 1 to 21, reflecting additional local
limit pushdown.
- Subquery/semi-join plans became more distributed:
    - TPC-DS CollectLeft hash joins: 615 → 610
    - TPC-DS partitioned hash joins: 98 → 103
    - TPC-DS left-semi occurrences: 11 → 25
    - TPC-DS network shuffles: 368 → 378
    - TPC-H - just a few
- These are meaningful topology changes: some subqueries now use
partitioned left-semi joins and therefore introduce hash shuffles
instead of collecting/broadcasting one side.
- Scalar rendering improved, especially decimal literals: internal forms
such as Some(0),7,2 now display as CAST(0.00 AS Decimal128(7, 2)).
- Minor changes (Ex. tpcds 21)
- `__common_expr_4` became `__common_expr_3`; that is only an internal
alias renumbering.
- The projection that renamed `d_date` to `__common_expr_2` disappeared.
- `d_date` is retained directly in the join output and referenced
directly by partial/final aggregates.
  - Column positions changed


- File-group allocation changed substantially
- Some explicit RoundRobinBatch repartitions disappeared and scans
gained different numbers of file groups
- Distribute byte ranges across partitions:
apache/datafusion#22439
- Lowers `repartition_file_min_size` from 10 MiB to 1 MiB. The PR
explicitly calls out TPC-DS SF1 dimension tables. Files may be
duplicated across multiple partitions where but each partition reads a
different byte range (this is hidden by <int>....<int>, but we know from
the correctness tests that nothing broke). A lot of tpcds queries now
split across `target_partitions` instead of staying under-partitioned.
In the `tpcds` plan tests, we use `target_partitions=3`.
Example:
```
-                │     t0: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t1: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t2: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t3: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t0: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t1: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t2: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t3: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
```

---------

Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com>
Co-authored-by: Gabriel <gabriel.musatmestre@datadoghq.com>
@jayshrivastava
jayshrivastava force-pushed the js/1-display-dynamic-filters branch from bac65ba to 75e43d6 Compare August 21, 2026 16:28
/// worry about any shared state.
///
/// [`update()`]: DynamicFilterPhysicalExpr::update()
pub(crate) fn sever_dynamic_filter_relationships_in_plan_for_display(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Here's one idea that is comes to mind:

What if we add yet another preparatory step for the plan in distributed_planner/, something like insert_broadcast or normalize_collect_joins, that severs all dynamic filter connections for good?

This would imply that dynamic filters will never be able to work through normal upstream mechanisms, and they should always be updated passing through the coordinator, even in the local case, but I do imagine this can simplify the overall approach, specially for future PRs.

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.

I would like to leave this for now.

The tricky part is that we want to sever some relationships but not all of them. For example, if there's local dynamic filters, a producer should be able to atomically update them in memory.

In future PRs, I detect local vs remote dynamic filters during static and dynamic planning, so we can think about severing some relationships then.

@jayshrivastava
jayshrivastava force-pushed the js/1-display-dynamic-filters branch from 131517e to 2f10b0b Compare September 2, 2026 20:48
Comment thread src/coordinator/distributed.rs Outdated
/// Execution state produced by distributed planning (static or dynamic) retained
/// for post-execution work such as plan rewrites to display metrics and dynamic filters.
#[derive(Debug, Clone)]
struct PreparedExecution {

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.

Now that TaskCtx is gone, we can now use PreparedPlan directly

let task_metrics = self.metrics_store.as_ref()?;
let plan = &self.prepared_plan.get()?.plan_for_viz;
Some(task_metrics.wait_for(&task_keys_for_plan(plan)).await)
}

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.

I decided this is the cleanest API: "wait until all task datas are present and then return them". It lets us remove the whole get() API on the store and just work with HashMap<...> directly

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

👍 yeap, this does look clean indeed, thanks!

}

/// Gathers metrics that belong to a task as a whole rather than to an execution-plan node.
fn stage_metrics(

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.

Basically a copy paste of gather_stage_header_metrics removed below

@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

@gabotechs Addressed the comments in the last few commits. I left the nontrivial comments open. I can check CI and rebase tomorrow 🫡

┌───── Stage 1 ── tasks=4, partitions=8
│ SortExec: TopK(fetch=5), expr=[id@0 ASC NULLS LAST], preserve_partitioning=[true]
│ DistributedLeafExec:
│ t0: DataSourceExec: file_groups={2 groups: [[/target/multi_task_collect_join_repros/build_side/part-0.parquet:<int>..<int>], [/target/multi_task_collect_join_repros/build_side/part-2.parquet:<int>..<int>]]}, projection=[id], file_type=parquet, predicate=DynamicFilter [ empty ], sort_order_for_reorder=[id@0 ASC NULLS LAST], dynamic_rg_pruning=eligible

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 lose these fields when doing the proto roundtrip before displaying. So, this means we actually lose these fields during execution since we serialize this plan before executing.

In a way, roundtripping the plan before displaying gives us a better picture of what's actually happening during execution.

I've addressed this here: apache/datafusion#24930, so these optimizations will come back in df56.

@gabotechs

Copy link
Copy Markdown
Collaborator

benchmarks run tpch/sf100

@gabot-0

gabot-0 commented Sep 4, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 417f80977d7f → PR head 7f006309f5b8 · View exact source diff

=== Comparing tpch/sf100 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=849.0, new=849.0, diff=no change (sum of per-query averages)
TOTAL: prev=82049 ms, new=81259 ms, diff=1.01 faster ✔
Show full query output
      q1: prev=3251 ms, new=2890 ms, diff=1.12 faster ✔, tasks: prev=18.0, new=18.0, diff=no change
      q2: prev=3485 ms, new=3668 ms, diff=1.05 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
      q3: prev=3397 ms, new=2968 ms, diff=1.14 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
      q4: prev=1441 ms, new=1437 ms, diff=1.00 faster ✔, tasks: prev=38.0, new=38.0, diff=no change
      q5: prev=4820 ms, new=4810 ms, diff=1.00 faster ✔, tasks: prev=52.0, new=52.0, diff=no change
      q6: prev=1510 ms, new=1758 ms, diff=1.16 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=10115 ms, new=10316 ms, diff=1.02 slower ✖, tasks: prev=53.0, new=53.0, diff=no change
      q8: prev=5518 ms, new=5374 ms, diff=1.03 faster ✔, tasks: prev=73.0, new=73.0, diff=no change
      q9: prev=9982 ms, new=10005 ms, diff=1.00 slower ✖, tasks: prev=77.0, new=77.0, diff=no change
     q10: prev=7865 ms, new=7664 ms, diff=1.03 faster ✔, tasks: prev=40.0, new=40.0, diff=no change
     q11: prev=3471 ms, new=3581 ms, diff=1.03 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q12: prev=2239 ms, new=2062 ms, diff=1.09 faster ✔, tasks: prev=44.0, new=44.0, diff=no change
     q13: prev=1991 ms, new=1959 ms, diff=1.02 faster ✔, tasks: prev=32.0, new=32.0, diff=no change
     q14: prev=1852 ms, new=1829 ms, diff=1.01 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q15: prev=4271 ms, new=4101 ms, diff=1.04 faster ✔, tasks: prev=30.0, new=30.0, diff=no change
     q16: prev= 826 ms, new= 972 ms, diff=1.18 slower ✖, tasks: prev=47.0, new=47.0, diff=no change
     q17: prev=5002 ms, new=5125 ms, diff=1.02 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q18: prev=5771 ms, new=5682 ms, diff=1.02 faster ✔, tasks: prev=68.0, new=68.0, diff=no change
     q19: prev=2179 ms, new=2209 ms, diff=1.01 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q20: prev=3063 ms, new=2849 ms, diff=1.08 faster ✔, tasks: prev=65.0, new=65.0, diff=no change
q21: Previously failed, and now also failed ❌
q22: Previously failed, and now also failed ❌
Verification and run details

Job 46 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-benchmarks --bin worker target from that checkout.

Identity PR base PR head
Source commit 417f80977d7fa5b6c4be0e4d1fa7f30864c5ca94 7f006309f5b8664ca838251a0ac0053559dbeb0c
Phase Base PR head
Build and deployment 1m 43s 1m 23s
All benchmarks 8m 39s 8m 28s
Benchmark tpch/sf100 8m 39s 8m 28s

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

Capacity: 12 c5n.2xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 20m 21s

@gabotechs gabotechs left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

💯 let's go!

@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

benchmarks run tpch/sf100

@gabot-0

gabot-0 commented Sep 4, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 417f80977d7f → PR head 7f006309f5b8 · View exact source diff

=== Comparing tpch/sf100 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=849.0, new=849.0, diff=no change (sum of per-query averages)
TOTAL: prev=80526 ms, new=81697 ms, diff=1.01 slower ✖
Show full query output
      q1: prev=2826 ms, new=2797 ms, diff=1.01 faster ✔, tasks: prev=18.0, new=18.0, diff=no change
      q2: prev=3287 ms, new=3455 ms, diff=1.05 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
      q3: prev=3139 ms, new=3044 ms, diff=1.03 faster ✔, tasks: prev=48.0, new=48.0, diff=no change
      q4: prev=1512 ms, new=1420 ms, diff=1.06 faster ✔, tasks: prev=38.0, new=38.0, diff=no change
      q5: prev=4618 ms, new=4973 ms, diff=1.08 slower ✖, tasks: prev=52.0, new=52.0, diff=no change
      q6: prev=1419 ms, new=1564 ms, diff=1.10 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=9901 ms, new=9935 ms, diff=1.00 slower ✖, tasks: prev=53.0, new=53.0, diff=no change
      q8: prev=5211 ms, new=5675 ms, diff=1.09 slower ✖, tasks: prev=73.0, new=73.0, diff=no change
      q9: prev=9713 ms, new=10151 ms, diff=1.05 slower ✖, tasks: prev=77.0, new=77.0, diff=no change
     q10: prev=7519 ms, new=8038 ms, diff=1.07 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q11: prev=3489 ms, new=3642 ms, diff=1.04 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q12: prev=2223 ms, new=2076 ms, diff=1.07 faster ✔, tasks: prev=44.0, new=44.0, diff=no change
     q13: prev=1917 ms, new=1946 ms, diff=1.02 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
     q14: prev=2150 ms, new=1783 ms, diff=1.21 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q15: prev=4503 ms, new=4086 ms, diff=1.10 faster ✔, tasks: prev=30.0, new=30.0, diff=no change
     q16: prev= 837 ms, new= 863 ms, diff=1.03 slower ✖, tasks: prev=47.0, new=47.0, diff=no change
     q17: prev=5190 ms, new=5221 ms, diff=1.01 slower ✖, tasks: prev=40.0, new=40.0, diff=no change
     q18: prev=5891 ms, new=5917 ms, diff=1.00 slower ✖, tasks: prev=68.0, new=68.0, diff=no change
     q19: prev=2130 ms, new=2126 ms, diff=1.00 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q20: prev=3051 ms, new=2985 ms, diff=1.02 faster ✔, tasks: prev=65.0, new=65.0, diff=no change
q21: Previously failed, and now also failed ❌
q22: Previously failed, and now also failed ❌
Verification and run details

Job 47 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-benchmarks --bin worker target from that checkout.

Identity PR base PR head
Source commit 417f80977d7fa5b6c4be0e4d1fa7f30864c5ca94 7f006309f5b8664ca838251a0ac0053559dbeb0c
Phase Base PR head
Build and deployment 2m 28s 1m 24s
All benchmarks 8m 28s 9m 2s
Benchmark tpch/sf100 8m 28s 9m 2s

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

Capacity: 12 c5n.2xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 21m 31s

@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

benchmarks run tpch/sf10

@gabot-0

gabot-0 commented Sep 4, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: PR base 417f80977d7f → PR head 7f006309f5b8 · View exact source diff

=== Comparing tpch/sf10 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=562.0, new=560.0, diff=2.0 fewer (0.4%) (sum of per-query averages)
TOTAL: prev=18866 ms, new=20422 ms, diff=1.08 slower ✖
Show full query output
      q1: prev= 344 ms, new= 630 ms, diff=1.83 slower ❌, tasks: prev=18.0, new=16.0, diff=2.0 fewer (11.1%)
      q2: prev= 816 ms, new= 887 ms, diff=1.09 slower ✖, tasks: prev=6.0, new=6.0, diff=no change
      q3: prev= 728 ms, new= 970 ms, diff=1.33 slower ✖, tasks: prev=28.0, new=28.0, diff=no change
      q4: prev= 326 ms, new= 586 ms, diff=1.80 slower ❌, tasks: prev=30.0, new=30.0, diff=no change
      q5: prev= 937 ms, new=1022 ms, diff=1.09 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
      q6: prev= 526 ms, new= 341 ms, diff=1.54 faster ✅, tasks: prev=12.0, new=12.0, diff=no change
      q7: prev=1294 ms, new=1346 ms, diff=1.04 slower ✖, tasks: prev=33.0, new=33.0, diff=no change
      q8: prev=1156 ms, new=1352 ms, diff=1.17 slower ✖, tasks: prev=54.0, new=54.0, diff=no change
      q9: prev=1510 ms, new=1595 ms, diff=1.06 slower ✖, tasks: prev=57.0, new=57.0, diff=no change
     q10: prev=1196 ms, new=1239 ms, diff=1.04 slower ✖, tasks: prev=20.0, new=20.0, diff=no change
     q11: prev= 587 ms, new= 560 ms, diff=1.05 faster ✔, tasks: prev=6.0, new=6.0, diff=no change
     q12: prev= 452 ms, new= 681 ms, diff=1.51 slower ❌, tasks: prev=30.0, new=30.0, diff=no change
     q13: prev= 592 ms, new= 627 ms, diff=1.06 slower ✖, tasks: prev=10.0, new=10.0, diff=no change
     q14: prev= 479 ms, new= 398 ms, diff=1.20 faster ✔, tasks: prev=21.0, new=21.0, diff=no change
     q15: prev=1074 ms, new=1035 ms, diff=1.04 faster ✔, tasks: prev=30.0, new=30.0, diff=no change
     q16: prev= 409 ms, new= 487 ms, diff=1.19 slower ✖, tasks: prev=3.0, new=3.0, diff=no change
     q17: prev= 769 ms, new= 743 ms, diff=1.03 faster ✔, tasks: prev=37.0, new=37.0, diff=no change
     q18: prev=1172 ms, new=1184 ms, diff=1.01 slower ✖, tasks: prev=45.0, new=45.0, diff=no change
     q19: prev= 427 ms, new= 569 ms, diff=1.33 slower ✖, tasks: prev=21.0, new=21.0, diff=no change
     q20: prev= 755 ms, new= 810 ms, diff=1.07 slower ✖, tasks: prev=32.0, new=32.0, diff=no change
     q21: prev=3020 ms, new=3054 ms, diff=1.01 slower ✖, tasks: prev=28.0, new=28.0, diff=no change
     q22: prev= 297 ms, new= 306 ms, diff=1.03 slower ✖, tasks: prev=9.0, new=9.0, diff=no change
Verification and run details

Job 48 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-benchmarks --bin worker target from that checkout.

Identity PR base PR head
Source commit 417f80977d7fa5b6c4be0e4d1fa7f30864c5ca94 7f006309f5b8664ca838251a0ac0053559dbeb0c
Phase Base PR head
Build and deployment 1m 22s 1m 22s
All benchmarks 2m 3s 2m 13s
Benchmark tpch/sf10 2m 3s 2m 13s

Workload: tpch/sf10 · all queries · 1 warmup + 5 measured iterations per query

Capacity: 12 c5n.2xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 7m 9s

Rich-T-kid pushed a commit to Rich-T-kid/datafusion-distributed that referenced this pull request Sep 5, 2026
## Summary 

Closes
datafusion-contrib#530

- This PR updates the upstream datafusion SHA to the HEAD of
https://github.com/apache/datafusion/commits/branch-55/ (edit: this
branch is continuously being updated. I will make sure this PR is at the
head before merging)
- Rust upgade to 1.94


## Changes

1. In `src/protobuf/distributed_codec.rs` we now use the
`proto_converter` argument during serde
- We still don't use the `DeduplicatingProtoConverter`, so dynamic
filters don't necessarily work. I think this is outside the scope of
this PR will be addressed in
datafusion-contrib#623,
which will be rebased after the upgrade.

3. `ExecutionPlan::apply_expressions` is added for every custom
`ExecutionPlan` in this repo
- Wrapper types (`MetricsWrapperExec`, `WorkUnitFileScanConfig`,
`DistributedLeafExec`) delegate to the inner type
- Other plans take`TreeNodeRecursion::Continue` because they have no
expressions (ex. `SamplerExec`)
- Note that `apply_expressions` does not need to yield sort or
partitioning expressions in the plan properties

3. We migrate from `partition_statistics` to `statistics_from_inputs`
for every `ExecutionPlan`.
- `src/distributed_planner/statistics/plan_statistics.rs` can just use
`statistics_from_inputs` directly instead of doing the
`StatisticsWrapper` workaround.

5. Range partitioning is now supported.
- CPU costing now includes range-key comparison cost and has a new unit
test. See src/distributed_planner/
   statistics/complexity_cpu.rs:238.
- I think there's open questions about range partitioning. I've opened
an issue here to make sure it behaves as expected after the upgrade:
datafusion-contrib#628 (comment)

6. Peak-memory metrics use the existing gauge wire representation.

DataFusion added MetricValue::PeakMemoryUsage. It is serialized as the
existing named-gauge protobuf variant to avoid a wire-format change. See
src/protocol/grpc/
   metrics_proto.rs:124.

The value and name survive, and aggregation is still additive, but
decoding produces a generic Gauge, not PeakMemoryUsage. The practical
difference is mainly display formatting: it
may render as a count rather than human-readable bytes. This is the
clearest remaining compromise/risk in the upgrade.

7. File-scan rebalancing changed its discriminator.

DataFusion removed partitioned_by_file_group;
output_partitioning.is_some() is now the source of truth. See
src/events/defaults/file_scan_config.rs:43. This decides whether files
are
round-robin rebalanced or split through FileGroupPartitioner, so it is
behavior-sensitive even though it is a one-line migration.

8. Two previously ignored correctness tests were enabled.
- See `tests/multi_task_collect_join_repros.rs`
- These were upstream DataFusion correctness fixes, not fixes made
locally in this upgrade.

9. drop(reporter) was made explicit on the sampler’s empty-input path.

The reporter sends its result on Drop; explicitly dropping it both
satisfies the new compiler/lint behavior and guarantees the zero-row EOS
report is sent before returning. See src/
   execution_plans/sampler.rs:259.

10. Plan changes

- `dynamic_rg_pruning=eligible` is now displayed on eligible scans:
1,354 occurrences in TPC-DS, 188 in TPC-H, and 12 in ClickBench
- `DataSourceExec` now displays its output partitioning. See
`tests/join.rs` (eventually, someone should delete this test
datafusion-contrib#628)
- Project after sort. This looks like some upstream optimizer rule
change ex. `tests/distributed_unions.rs` and
`tests/distributed_aggregation.rs`.
```
-          │   SortExec: expr=[MinTemp@0 ASC NULLS LAST, RainToday@1 ASC NULLS LAST], preserve_partitioning=[true]
-          │     ProjectionExec: expr=[MaxTemp@0 as MinTemp, RainToday@1 as RainToday]
+          │   ProjectionExec: expr=[MaxTemp@0 as MinTemp, RainToday@1 as RainToday]
+          │     SortExec: expr=[MaxTemp@0 ASC NULLS LAST, RainToday@1 ASC NULLS LAST], preserve_partitioning=[true]
```
- LocalLimitExec became more common: TPC-DS went from 0 to 20
occurrences and ClickBench from 1 to 21, reflecting additional local
limit pushdown.
- Subquery/semi-join plans became more distributed:
    - TPC-DS CollectLeft hash joins: 615 → 610
    - TPC-DS partitioned hash joins: 98 → 103
    - TPC-DS left-semi occurrences: 11 → 25
    - TPC-DS network shuffles: 368 → 378
    - TPC-H - just a few
- These are meaningful topology changes: some subqueries now use
partitioned left-semi joins and therefore introduce hash shuffles
instead of collecting/broadcasting one side.
- Scalar rendering improved, especially decimal literals: internal forms
such as Some(0),7,2 now display as CAST(0.00 AS Decimal128(7, 2)).
- Minor changes (Ex. tpcds 21)
- `__common_expr_4` became `__common_expr_3`; that is only an internal
alias renumbering.
- The projection that renamed `d_date` to `__common_expr_2` disappeared.
- `d_date` is retained directly in the join output and referenced
directly by partial/final aggregates.
  - Column positions changed


- File-group allocation changed substantially
- Some explicit RoundRobinBatch repartitions disappeared and scans
gained different numbers of file groups
- Distribute byte ranges across partitions:
apache/datafusion#22439
- Lowers `repartition_file_min_size` from 10 MiB to 1 MiB. The PR
explicitly calls out TPC-DS SF1 dimension tables. Files may be
duplicated across multiple partitions where but each partition reads a
different byte range (this is hidden by <int>....<int>, but we know from
the correctness tests that nothing broke). A lot of tpcds queries now
split across `target_partitions` instead of staying under-partitioned.
In the `tpcds` plan tests, we use `target_partitions=3`.
Example:
```
-                │     t0: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t1: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t2: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
-                │     t3: DataSourceExec: file_groups={2 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t0: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t1: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t2: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
+                │     t3: DataSourceExec: file_groups={3 groups: [[/testdata/tpcds/plans_sf1_partitions4/date_dim/part-0.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-1.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>, /testdata/tpcds/plans_sf1_partitions4/date_dim/part-2.parquet:<int>..<int>], [/testdata/tpcds/plans_sf1_partitions4/date_dim/part-3.parquet:<int>..<int>]]}, projection=[d_date_sk, d_week_seq, d_day_name], file_type=parquet, predicate=DynamicFilter [ empty ]
```

---------

Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com>
Co-authored-by: Gabriel <gabriel.musatmestre@datadoghq.com>
@jayshrivastava
jayshrivastava merged commit e5bca18 into main Sep 8, 2026
63 of 64 checks passed
@jayshrivastava
jayshrivastava deleted the js/1-display-dynamic-filters branch September 8, 2026 12:56
jayshrivastava added a commit that referenced this pull request Sep 21, 2026
## Stack

This stack of PRs implements distributed dynamic filtering #528 
1. #623
2. #634
<- you are here
3. #635
4. #636
5. #637
6. #639

## Goal

The coordinator should know what dynamic filters exist and where to
route updates.

## Details

### 1. Dynamic Filter Registry

```
  QueryCoordinator
  └── DynamicFilterRegistry
      └── filters: Map<expression_id, PlannedDynamicFilter>
```

Each `PlannedDynamicFilter` stores
- the producers and their stage/tasks
- the consumers and their stage/tasks

This will be used in future PRs to store incoming dynamic filter updates
from workers and determine how/where to forward the updates.

#### Implementation

In the `StageCoordinator`, we send every task to the registry and
extract dynamic filters.

### 2. Network Anchors

In this situation, the hash join does an `execute()`-time check to
determine if it should update its dynamic filter. It checks to see if
the filter is used by any children using `apply_expressions` (before
`apply_expressions` was added upstream, it was an Arc pointer strong
count check to see if there were multiple references).
```
worker 1
HashJoinExec  (dynamic_filter_predicate)
    NetworkShuffleExec

worker 2
DataSourceExec (dynamic_filter_predicate)
```
The join sees that no plan nodes below it use the filter, so it decides
not to update it.

Ideally, the hash join decides at optimization time, before distributed
planning. I've opened a discussion here about it:
apache/datafusion#18856 (comment).
While that issue is being resolved, I propose this workaround:

We create an "anchor" to make it seem like the `NetworkShuffleExec` uses
the filter.
```
worker 1
HashJoinExec  (dynamic_filter_predicate)
    NetworkShuffleExec  (anchor: dynamic_filter_predicate)

worker 2
DataSourceExec (dynamic_filter_predicate)
```

#### Network Anchors Implementation

The implementation adds serialization overhead but is simpler. In static
and dynamic planning, we recursively propagate all anchors upwards in
the plan to all the network boundaries. We can revisit this
implementation in future iterations. This recursive implementation is in
`inject_network_boundaries`.
```
stage3:
    HashJoinExec <- producer of filter1
        NetworkShuffleExec  (anchors: filter1, filter2)

stage2:
    RepartitionExec
        AggregateExec <- producer of filter2
            NetworkShuffleExec   (anchors: filter1, filter2)

stage1:
        DataSourceExec (consumer: filter1, filter2)
```
This means we serialize 8 filters in total.

However, the minimal anchors you need are like this:
```
stage3:
    HashJoinExec <- producer #1
        NetworkShuffleExec  (anchors: filter2)

stage2:
    RepartitionExec
        AggregateExec <- producer #2
            NetworkShuffleExec   (anchors: filter2)

stage1:
        DataSourceExec (consumer: filter1, filter2)
```
In this plan, we would serialize 6 filters.

For 1 dynamic filter, the minimum filters you need to serialize are 1
(producer) + N (consumers) + 1 (network boundary). In this
implementation, we serialize 1 (producer) + N (consumers) + M (all
network boundaries above the consumer)

### Other Notes

See
#528.
During dynamic planning, the sampler on the probe side of a hash join
may overreport rows / cost because dynamic filters aren't being applied
yet.
jayshrivastava added a commit that referenced this pull request Sep 21, 2026
## Stack

This stack of PRs implements distributed dynamic filtering #528 
1. #623
2. #634
3. #635
<- you are here
4. #636
5. #637
6. #639


## Problem

The `QueryCoordinator` needs to receive partial dynamic filter updates
from workers.

## Solution
We introduce a new `WorkerToCoordinatorMsg` which 
```
message ProducedDynamicFilter {
  uint64 expression_id = 1;
  // Serialized datafusion.proto.PhysicalExprNode.
  bytes expression_proto = 2;
}
```

In this PR makes each worker unconditionally send updates (via
`wait_update()` and `wait_complete()`) to the coordinator for any
`dynamic_filter_remote_producer_ids` in the `SetPlanRequest`. The
purpose of `dynamic_filter_remote_producer_ids` is to exclude any
dynamic filters who only have local consumers - these don't need to be
forwarded to the coordinator.
@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

benchmarks run clickbench/0-100-date32 --base e5a4a77

@gabot-0

gabot-0 commented Sep 22, 2026

Copy link
Copy Markdown

@jayshrivastava Invalid base e5a4a779d191206c1acd5cbfb31f3fed87c57934; only main is supported.

How to use the benchmark bot

Post a comment whose first non-empty line is:

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

For example:

benchmarks run tpch/sf10 tpch/sf100 --instance-type m5.2xlarge --nodes 24 --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.
--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 and iteration overrides are not supported; every query uses 1 warmup and 5 measured iterations for both revisions.

jayshrivastava added a commit that referenced this pull request Sep 25, 2026
## Stack

This stack of PRs implements distributed dynamic filtering #528 
1. #623
2. #634
3. #635
4. #636
<- you are here
5. #637
6. #639

## Details

This PR adds machinery around merging dynamic filters in the
`DynamicFilterRegistry`.

```
    register_task(stage1, task1)                                                                                              
    register_task(stage1, task2)       ───────────┐                                                                              
    register_task(stage1, task3)                  │          ┌───────────────────────┐      merge() when                         
         seal_stage(stage1)                       ├─────────▶│ DynamicFilterRegistry │───▶  - stage is sealed; and               
                                                  │          └───────────────────────┘      - there's enough partial filter      
                                                  │                                         updates                              
record_dynamic_filter_update(stage1, task1)   ────┘                                                                              
record_dynamic_filter_update(stage2, task1)                                                                                           
record_dynamic_filter_update(stage3, task1)                                                                                           
                                                                                                                              
```

The query coordinator calls `register_task` for each task in a stage.
Concurrently, any running task from any stage can send a dynamic filter
update to the registry via `record_dynamic_filter_update`. The registry
needs to detect when all the updates are present and `merge()` the
partial dynamic filters. To help detect this, the coordinator is
responsible for calling `seal_stage` once all the tasks have commenced
so we know that no tasks will be added in the future.

The PR implements the above. In the next 2 PRs, we will actually forward
the merged filters to consumers.
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.

[dynamic filtering] 2. collect and display dynamic filters in plans

4 participants