coordinator: merge partial dynamic filters - #636
Conversation
2cc815b to
7f6de9a
Compare
400db1c to
ad3ea34
Compare
b798e6c to
876ac94
Compare
876ac94 to
d83e14d
Compare
d83e14d to
593f037
Compare
593f037 to
1548203
Compare
1548203 to
21b4e38
Compare
21b4e38 to
4e36655
Compare
4e36655 to
441ba38
Compare
8af79a1 to
40c52ae
Compare
40c52ae to
5528a7a
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 <- 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.
75c7432 to
172cecb
Compare
3f6f3d6 to
1004f74
Compare
gabotechs
left a comment
There was a problem hiding this comment.
Looking good! left some small comments, nothing structural.
| /// Latest accepted snapshot from each producer task. | ||
| pub(super) producer_filters: HashMap<TaskKey, PhysicalDynamicFilterNode>, | ||
| /// Full dynamic filter containing the merged predicate and its completion state. | ||
| pub(super) merged: Option<PhysicalDynamicFilterNode>, |
There was a problem hiding this comment.
I see this is working in terms of protobuf messages rather than raw DataFusion expressions.
This means that proto conversion would happen even in a fully in-memory setup. Is it easy to avoid and just work with normal DataFusion types here rather than protobuf messages?
There was a problem hiding this comment.
The main reason is honestly that is_complete isn't public on DynamicFilterPhysicalExpr. So there's actually no way to know if it's complete or not here.
We could detect it on the worker here and send a is_complete bit in the `ProducedDynamicFilter. Do you think making that change is worth it?
There was a problem hiding this comment.
🤔 I don't think there should be any reason for is_complete() to be private upstream, even the mark_complete() method is public.
Just put a PR for that upstream (apache/datafusion#25738). In the meantime, it's fine to keep what you have here 👍
Merge the latest TopK and MIN/MAX producer bounds incrementally. Keep complete-producer coverage for partitioned joins and the complete-replica fast path for CollectLeft joins. Preserve the full dynamic filter and distinguish update readiness from global completion.
… before the limit
cd6fe19 to
52342b3
Compare
Stack
This stack of PRs implements distributed dynamic filtering #528
Details
This PR adds machinery around merging dynamic filters in the
DynamicFilterRegistry.The query coordinator calls
register_taskfor each task in a stage. Concurrently, any running task from any stage can send a dynamic filter update to the registry viarecord_dynamic_filter_update. The registry needs to detect when all the updates are present andmerge()the partial dynamic filters. To help detect this, the coordinator is responsible for callingseal_stageonce 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.