Skip to content

coordinator: forward merged dynamic filters to consumers - #637

Open
jayshrivastava wants to merge 5 commits into
mainfrom
js/5-forward-dynamic-filters
Open

jayshrivastava wants to merge 5 commits into
mainfrom
js/5-forward-dynamic-filters

Conversation

@jayshrivastava

@jayshrivastava jayshrivastava commented Aug 13, 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
  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 <- you are here
  6. worker: apply merged dynamic filters during execution #639

Details

  • Route each merged dynamic filter to the tasks that consume them using the existing CoordinatorToWorker channel
  • Refactor the CoordinatorToWorker lifetime handling to allow async senders to propagate errors back to the coordinator (ie. if sending a merged dynamic filter fails). We know have a global CancellationToken that fires when the query ends or when there is an error.

@jayshrivastava jayshrivastava changed the title feat: forward merged dynamic filters to consumers coordinator: forward merged dynamic filters to consumers Aug 13, 2026
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 3229903 to ad77338 Compare August 13, 2026 19:27
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from ad77338 to 15e6602 Compare August 17, 2026 18:54
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 2 times, most recently from 56fee53 to 1572836 Compare August 17, 2026 19:38
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 1572836 to 296026d Compare August 18, 2026 18:15
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 2 times, most recently from e04ceac to e720301 Compare August 21, 2026 16:28
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from e720301 to 7a74a70 Compare August 21, 2026 16:48
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 7a74a70 to c1b6398 Compare August 21, 2026 20:58
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from c1b6398 to 6a70c73 Compare August 22, 2026 15:00
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 2 times, most recently from 0a68111 to 42cb2ac Compare August 23, 2026 15:27
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 6e97fe7 to 12d9331 Compare September 21, 2026 20:15
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 2 times, most recently from 33be498 to 5863fe4 Compare September 21, 2026 21:28
@jayshrivastava jayshrivastava changed the title [do not review] coordinator: forward merged dynamic filters to consumers coordinator: forward merged dynamic filters to consumers Sep 21, 2026
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 3 times, most recently from acdb19a to da1fa7f Compare September 22, 2026 03:55
Comment thread src/coordinator/distributed.rs Outdated
Some(error) => Err(error),
None => Ok(()),
}
});

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 error channel is a bit controversial. It adds a way for us to propagate errors to the main query.

Comment thread src/coordinator/query_coordinator.rs Outdated
let _ = error_tx.try_send(error);
break;
}
};

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 debatable.

Previously, errors in the worker -> coordinator channel were reported in the joinset, after the query ran. We have some important messages now

  • ProducedDynamicFilter (in vanilla datafusion, dynamic filtering errors are not silently discarded, so I imagine we would want these to fail loudly)
  • LoadInfo, important for dynamic planning
  • In this PR, I also report coordinator -> worker request serialization errors back on this channel (see below)

So, I figured we should propagate this error to the query coordinator.

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.

👍 I see the point in propagating errors. I think we might have better options that yet another side-channel for this.

In distributed.rs, we build one RecordBatchReceiverStreamBuilder that eventually gets converted to the output stream. Under the hood, RecordBatchReceiverStreamBuilder has a JoinSet that already acts as a side channel listening for errors happening in any of the spawned tasks, surfacing an error immediately after appearing.

Note how we also have another join_set as part of the QueryCoordinator struct, but errors happening there are not really surfaced anywhere (unlike the join_set from RecordBatchReceiverStreamBuilder).

I think there's an opportunity here to consolidate things, some options that come to mind are:

  • Just extract the logic we care about from RecordBatchReceiverStreamBuilder and adapt it to QueryCoordinator so that errors in spawned tasks from QueryCoordinator.join_set are automatically surfaced without side channels.
  • Bake the RecordBatchReceiverStreamBuilder inside QueryCoordinator and just spawn everything it's join_set, removing ours.

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 like the idea of using a joinset. Will look into this.

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.

Just pushed a new change using a joinset! I think it's pretty clean :)

Comment thread src/coordinator/query_coordinator.rs Outdated
let msg = match msg {
Ok(msg) => msg,
Err(error) => {
let _ = error_tx.try_send(error);

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'm not sure if this is clear enough. The error_tx has a buffer size of 1, so we can ignore errors due to the channel being full. Is it worth adding a utility for this?

let plan_bytes_sent = set_plan_request.plan_proto.len();
let task_ctx = Arc::clone(ctx);
// Tonic request streams cannot yield errors, so return encoding failures through
// the response stream instead. Only the first failure matters.

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 what I was referencing in my comment above. I decided to report request errors on the worker -> coordinator channel - can't think of a cleaner way at the moment.

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.

🤔 yeah, I've been thinking a bit, and I also cannot figure out a way of propagating these errors without a side channel...

Comment thread src/coordinator/query_coordinator.rs Outdated
// worker->coordinator channel, and then that channel is closed.
.chain(keep_stream_alive(Arc::clone(self.end_stream_notifier)))
// Keep the channel open after work-unit delivery until the query finishes.
.take_until(self.query_finished.clone().cancelled_owned())

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.

The old pattern didn't work for the senders in the dynamic filter registry.

The old notifier required waiting for everything before it to finish, but the tx's in the dynamic filter registry don't necessarily finish, they stay open. So, I opened for a shared cancellation token which will cancel the stream even if the senders in the registry are still alive.

@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from da1fa7f to 8e449ba Compare September 22, 2026 14:46
@jayshrivastava

Copy link
Copy Markdown
Collaborator Author

benchmarks run clickbench/0-100-date32 --iterations 20 --base main

@gabot-0

gabot-0 commented Sep 22, 2026 •

Copy link
Copy Markdown

Requested by this comment.

Benchmark results

Compared: Main 0f6fe65e5a5f → PR head 8e449ba3a75c · View exact source diff

=== Comparing clickbench/0-100-date32 results 'datafusion-benchmark-base' [prev] with 'datafusion-benchmark-head' [new] ===
TASKS: prev=912.0, new=912.0, diff=no change (sum of per-query averages)
TOTAL: prev=72369 ms, new=68665 ms, diff=1.05 faster ✔
Show full query output
      q0: prev=   2 ms, new=   1 ms, diff=2.00 faster ✅, tasks: prev=0.0, new=0.0, diff=no change
      q1: prev= 405 ms, new= 360 ms, diff=1.13 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q2: prev= 292 ms, new= 242 ms, diff=1.21 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q3: prev= 398 ms, new= 358 ms, diff=1.11 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
      q4: prev= 496 ms, new= 451 ms, diff=1.10 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
      q5: prev= 686 ms, new= 679 ms, diff=1.01 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
      q6: prev=   2 ms, new=   2 ms, diff=1.00 slower ✖, tasks: prev=0.0, new=0.0, diff=no change
      q7: prev= 420 ms, new= 406 ms, diff=1.03 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
      q8: prev= 547 ms, new= 490 ms, diff=1.12 faster ✔, tasks: prev=36.0, new=36.0, diff=no change
      q9: prev= 619 ms, new= 495 ms, diff=1.25 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q10: prev= 898 ms, new= 540 ms, diff=1.66 faster ✅, tasks: prev=36.0, new=36.0, diff=no change
     q11: prev= 752 ms, new= 499 ms, diff=1.51 faster ✅, tasks: prev=36.0, new=36.0, diff=no change
     q12: prev= 843 ms, new= 747 ms, diff=1.13 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q13: prev= 893 ms, new= 885 ms, diff=1.01 faster ✔, tasks: prev=36.0, new=36.0, diff=no change
     q14: prev= 831 ms, new= 772 ms, diff=1.08 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q15: prev= 661 ms, new= 361 ms, diff=1.83 faster ✅, tasks: prev=24.0, new=24.0, diff=no change
     q16: prev= 800 ms, new= 755 ms, diff=1.06 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q17: prev= 801 ms, new= 737 ms, diff=1.09 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q18: prev= 977 ms, new= 763 ms, diff=1.28 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q19: prev= 473 ms, new= 459 ms, diff=1.03 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
     q20: prev=6844 ms, new=5981 ms, diff=1.14 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
     q21: prev=5777 ms, new=4797 ms, diff=1.20 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q22: prev=5164 ms, new=4768 ms, diff=1.08 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q23: prev=16294 ms, new=15296 ms, diff=1.07 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
     q24: prev= 552 ms, new= 534 ms, diff=1.03 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
     q25: prev= 586 ms, new= 574 ms, diff=1.02 faster ✔, tasks: prev=12.0, new=12.0, diff=no change
     q26: prev= 550 ms, new= 578 ms, diff=1.05 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
     q27: prev=5520 ms, new=5171 ms, diff=1.07 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q28: prev=2201 ms, new=2492 ms, diff=1.13 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q29: prev= 257 ms, new= 298 ms, diff=1.16 slower ✖, tasks: prev=12.0, new=12.0, diff=no change
     q30: prev= 735 ms, new= 774 ms, diff=1.05 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q31: prev= 730 ms, new= 764 ms, diff=1.05 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q32: prev= 655 ms, new= 639 ms, diff=1.03 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q33: prev=5255 ms, new=5603 ms, diff=1.07 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q34: prev=5269 ms, new=5888 ms, diff=1.12 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q35: prev= 379 ms, new= 492 ms, diff=1.30 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q36: prev= 772 ms, new= 774 ms, diff=1.00 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q37: prev= 535 ms, new= 527 ms, diff=1.02 faster ✔, tasks: prev=24.0, new=24.0, diff=no change
     q38: prev= 701 ms, new= 747 ms, diff=1.07 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q39: prev= 884 ms, new= 909 ms, diff=1.03 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q40: prev= 298 ms, new= 301 ms, diff=1.01 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q41: prev= 306 ms, new= 352 ms, diff=1.15 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
     q42: prev= 309 ms, new= 404 ms, diff=1.31 slower ✖, tasks: prev=24.0, new=24.0, diff=no change
Verification and run details

Job 184 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 Main PR head
Source commit 0f6fe65e5a5fd305e63dfc535b47a643b7cda49e 8e449ba3a75cac1b9b86432dab3e99b5655aba4d
Phase Base PR head
Build and deployment 1m 51s 1m 36s
All benchmarks 27m 58s 27m 18s
Benchmark clickbench/0-100-date32 27m 58s 27m 18s

Workload: clickbench/0-100-date32 · all queries · 1 warmup + 20 measured iterations per query for both revisions

Capacity: 12 c5n.4xlarge nodes for both revisions

Other timings: Queue 1s · Dataset validation 0s · Total 58m 53s

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.

@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 8e449ba to abd172e Compare September 23, 2026 01:23
Comment thread src/coordinator/query_coordinator.rs Outdated
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch 2 times, most recently from c5a469d to 2f0a742 Compare September 24, 2026 16:24
Comment thread src/coordinator/query_task_state.rs Outdated
// Surface errors in any tasks.
let tasks_finished = loop {
let result = {
let mut tasks = self.task_state.join_set.lock().unwrap();

@gabotechs gabotechs Sep 25, 2026 •

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.

This is a bit scary.... Acquiring a lock in every poll sounds excessive.

Having a lock-free solution here should be easy: just rely on channels and tokio::select! for that. I'd take some inspiration on ReceiverStreamBuilder upstream instead, there it's handled very well.

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.

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.

Moved to use a select! based approach. Left some comments in query_task_state.rs as well.

Comment thread tests/work_unit_feed.rs
Comment on lines +720 to +723
let query = r#"
SELECT * FROM test_work_unit('a', 2, 'wait(30000), rows(1)', 'rows(1), err(boom_distributed)')
"#;
// Expect the query to fail fast rather than waitinng for the unrelated 30s feed.

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.

Why did this test needed to change here?

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.

Previously, we asserted that the error propagates. Now we assert it propagates fast. I figured it's a small diff and provides some value.

Base automatically changed from js/4-merge-dynamic-filters to main September 25, 2026 14:52
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.
Carry full dynamic filters through MaybeEncoded without exposing protobuf types in the transport API. Encode each merged snapshot once, enqueue updates under the registry lock, and send the latest state to late consumers. Keep delivery fail-open and release senders when the query ends.
@jayshrivastava
jayshrivastava force-pushed the js/5-forward-dynamic-filters branch from 2f0a742 to fc39b13 Compare September 28, 2026 15:39

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