feat: forward remote dynamic filter updates to coordinator - #635
Conversation
60dc39a to
7bee97e
Compare
7bee97e to
53be7d0
Compare
4fd1433 to
c29bcf1
Compare
c29bcf1 to
a7bb152
Compare
a7bb152 to
8bc2dc9
Compare
e1a6b47 to
77c9528
Compare
77c9528 to
dd39b56
Compare
dd39b56 to
fceebde
Compare
5912188 to
a60f964
Compare
323cbf2 to
336d392
Compare
336d392 to
8a3f76f
Compare
| } | ||
| // Runtime dynamic-filter reports are accepted by this transport change. A | ||
| // later change in the stack will retain and merge them. | ||
| WorkerToCoordinatorMsg::ProducedDynamicFilter(_) => {} |
There was a problem hiding this comment.
In this PR, we do nothing when receiving dynamic filter updates. In later PRs, we will consume these updates.
|
benchmarks run tpch/sf100 |
|
Requested by this comment. Benchmark job 51 failed for |
| let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); | ||
|
|
||
| let local_filter = dynamic_filter(); | ||
| let local_probe = Arc::new(FilterExec::try_new(local_filter.clone(), empty(&schema))?) | ||
| as Arc<dyn ExecutionPlan>; | ||
| let local_join = join_with_filter(&schema, local_probe, Arc::clone(&local_filter))?; | ||
| assert!(dynamic_filter_remote_producer_ids(&local_join)?.is_empty()); | ||
|
|
||
| let remote_filter = dynamic_filter(); | ||
| let repartition = Arc::new(RepartitionExec::try_new( | ||
| empty(&schema), | ||
| Partitioning::Hash(vec![Arc::new(Column::new("a", 0))], 1), | ||
| )?) as Arc<dyn ExecutionPlan>; | ||
| let remote_probe = Arc::new( | ||
| NetworkShuffleExec::try_new(repartition, 1)? | ||
| .with_dynamic_filter_anchors(vec![remote_filter.clone()]), | ||
| ) as Arc<dyn ExecutionPlan>; |
There was a problem hiding this comment.
I see this being a pattern with the tests: plans get constructed manually by chaining operators together.
It's pretty hard to see what's happening and what's getting chained with what to a human eye. Do you think there's a chance we can just use normal SQL for these tests?
The main reason for having the weather and flights datasets committed to the codebase is so that we can just use SQL for building plans, rather than constructing them manually.
There was a problem hiding this comment.
Done. I migrated the function dynamic_filter_remote_producer_ids to src/dynamic_filtering/discovery.rs and we now test it in the snapshots in tests/dynamic_filtering/discovery.rs
| let expression = Arc::clone(&dynamic_filter) as Arc<dyn PhysicalExpr>; | ||
| let (_cancel_tx, cancel_rx) = watch::channel(false); | ||
| let task_ctx = SessionContext::new().task_ctx(); | ||
| let mut stream = produced_dynamic_filter_stream(expression_id, expression, cancel_rx); |
There was a problem hiding this comment.
I really think we should change the approach to testing in general in this PR.
The produced_dynamic_filter_stream is a private function of this module, and therefore an implementation detail subject to change. By testing at this level, we are testing the implementation detail rather than the overall functionality, and the tests tend to be very verbose.
Here's one idea for change the way we approach testing in this PR: for anything related to dynamic filter collection and over-the-wire transfer, we can add some nice DataFusion metrics at the coordinator level, and during integration testing, we can perform assertion on those metrics.
With that, we'd win two things:
- Reliable tests that survive changes to the implementation details.
- Runtime metrics that give visibility around what happened.
There was a problem hiding this comment.
I added a metric to track how many updates are received. I think this is good enough for now, we just assert that they are nonzero. Eventually, once remote dynamic filtering works, we will assert that the filters show up in plans and that the results match with dyn filtering on vs off.
22137b1 to
8337f49
Compare
8337f49 to
d579734
Compare
dcbb41c to
5d04683
Compare
|
This is ready for another review. |
gabotechs
left a comment
There was a problem hiding this comment.
👍 nice! nothing major
| use datafusion::physical_expr::expressions::{Column, lit}; | ||
|
|
||
| #[tokio::test] | ||
| async fn cancellation_stops_dynamic_filter_updates() { |
There was a problem hiding this comment.
This test is essentially just testing that tokio::select! inside a private function works. This is one of type of tests that is subject to change as implementation details change.
A more complete coverage could be to run a relatively long query (e.g. TPCH Q1), abort it mid-way, and make sure the query ends immediately, without dyn filters getting any update.
There was a problem hiding this comment.
I guess what we want is to assert the streams all close (which we don't today). In tests/stateful_data_cleanup.rs, we only test that the task data entries are evicted. I'm going to add a little test-only metric that let's us observe streams getting cleaned up.
| sql: &str, | ||
| expected_rows: usize, | ||
| collect_dynamic_filters: bool, | ||
| expect_dynamic_filter_updates: Option<bool>, |
There was a problem hiding this comment.
Why would this function accept None instead of just a false?
Task-local consumer
producer F ────────────────> consumer F
shared expression; no RPC
Cross-task consumer
producer F ── boundary anchor F ──> SetPlanRequest.report_ids=[F]
│
├── wait_update() -----+
├── wait_complete() ---+--> latest observed F --RPC--> coordinator
└── query cancellation +--> stop
Derive the report allowlist from producer IDs intersected with network-boundary anchor IDs. Workers observe only allowlisted producers, while DataFusion continues updating task-local consumers directly in memory. Watch-based updates may naturally coalesce, and completion remains observed separately because it does not advance the generation.
45ecf27 to
c071afb
Compare
## 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.
## 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.
## What Bump the datafusion-distributed ref so we stay current. It picks up distributed dynamic filters (datafusion-contrib/datafusion-distributed#634, datafusion-contrib/datafusion-distributed#635), reduced producer fan-out in NetworkShuffleExec (datafusion-contrib/datafusion-distributed#717), the AQE worker-assignment retry (datafusion-contrib/datafusion-distributed#723), fractional desired task counts (datafusion-contrib/datafusion-distributed#714), and the datafusion-iceberg refactor (datafusion-contrib/datafusion-distributed#742). The introduction of distributed dynamic filters introduces a large hash join regression when used over our transport layer, so this PR disables their use until we get out a fix. ## Tests CI's green. Once benchmarks run we should be able to verify this _does not_ introduce the aforementioned perf regression. --------- Co-authored-by: paradedb-github-bot[bot] <282009505+paradedb-github-bot[bot]@users.noreply.github.com>
Stack
This stack of PRs implements distributed dynamic filtering #528
Problem
The
QueryCoordinatorneeds to receive partial dynamic filter updates from workers.Solution
We introduce a new
WorkerToCoordinatorMsgwhichIn this PR makes each worker unconditionally send updates (via
wait_update()andwait_complete()) to the coordinator for anydynamic_filter_remote_producer_idsin theSetPlanRequest. The purpose ofdynamic_filter_remote_producer_idsis to exclude any dynamic filters who only have local consumers - these don't need to be forwarded to the coordinator.