Is your feature request related to a problem or challenge?
I think this would help solve two problems
- It allows for certain optimizations
- Solves a problem related to distributed datafusion (and perhaps other frameworks like ballista + comet could benefit).
1) Remove CASE hash(expr)
Related: #23376 (comment)
To support cases like below where the partitioning of the dynamic filter producer (hash join) is different than the consumer (data source), the hash join uses a CASE hash(expr) to "bake" partitioning in to the join.
DynamicFilterPhysicalExpr:
┌────────────────────────┐
│ HashJoinExec: │ CASE Hash(a@0) % 4
│mode=Partitioned on=a@0 │ WHEN 0: a@0 >= v0 AND a@0 <= v1
└────────────────────────┘ WHEN 1: a@0 IN (SET) ([v2, v3, v4 ...])
▲ ▲ ... (1 case per partition)
│ │
┌────────────────┘ └────────────────┐ │
│ │ Pushed down to data source
┌──────────────────────────┐ ┌──────────────────────────┐
│ DataSourceExec: │ │ RepartitionExec: │ │
│Partitioning=Hash(a, 4) │ │ Partitioning=Hash(a, 4) │ │
└──────────────────────────┘ └──────────────────────────┘ │
│ │
│ │
│ │
│ │
... any number of RepartitionExecs or │
plan nodes. It does not matter │
│ │
│ │
│ ▼
┌──────────────────────────┐
│ DataSourceExec: │ Each partition uses the same DynamicFilterPhysicalExpr, except
│ Partitioning=Unknown(10) │ each row will only hash to one case.
└──────────────────────────┘
What if the partitioning of the data source is the same as the hash join? In this case, the CASE hash(expr) is pure overhead.
┌────────────────────────┐
│ HashJoinExec: │
│mode=CollectLeft on=a@0 │
└────────────────────────┘
▲ ▲
│ │
┌────────────────┘ └────────────────┐
│ │
┌──────────────────────────┐ ┌──────────────────────────┐
│ DataSourceExec: │ │ RepartitionExec: │
│Partitioning=Hash(a, 4) │ │ Partitioning=Hash(a, 4) │
└──────────────────────────┘ └──────────────────────────┘
▲
│
┌──────────────────────────┐
│ DataSourceExec: │
│ Partitioning=Hash(4) │
└──────────────────────────┘
It would be nice if the DataSourceExec could consume per-partition dynamic filters. I think most DataSourceExec have to downcast PhysicalExpr to DynamicFilterPhysicalExpr to consume them anyways. If the DynamicFilterPhysicalExpr stored expressions per partition, the DataSourceExec could apply per-partition filters rather than use the CASE hash(expr)
2) Distributed Execution
More tails in the RFC: datafusion-contrib/datafusion-distributed#553
In datafusion-distributed, I'm interested in implementing dynamic filtering similarly to Trino.
- During execution, the dynamic filter producers (
HashJoinExec, SortExec, AggregateExec) will send dynamic filter updates to the coordinator.
- The coordinator will union / merge the dynamic filter updates from all workers.
- The coordinator will send the unioned / merged dynamic filter updates to the consumers (
DataSourceExec).
For correctness, the coordinator must wait for all dynamic filter updates from all workers before sending them to the consumers, otherwise, the consumers may prune rows incorrectly.
In the diagram below, each task in stage 2 creates a dynamic filter, so there's 3 in total:
Task 1:
CASE Hash(a@0) % 4
WHEN 0: a@0 >= v0 AND a@0 <= v1
... 4 cases in total
Task 2:
CASE Hash(a@0) % 4
WHEN 0: a@0 IN LIST [v2, v3, v4 ...]
... 4 cases in total
Task 3:
CASE Hash(a@0) % 4
WHEN 0: a@0 >= v6
... 4 cases in total
The plan is to merge them together.
Option 1: OR them
CASE Hash(a@0) % 4
WHEN 0: a@0 >= v0 AND a@0 <= v1
... 4 cases in total
OR
CASE Hash(a@0) % 4
WHEN 0: a@0 IN LIST [v2, v3, v4 ...]
... 4 cases in total
OR
CASE Hash(a@0) % 4
WHEN 0: a@0 >= v6
... 4 cases in total
Option 2: Case-wise OR
CASE Hash(a@0)
WHEN 0: a@0 >= v0 AND a@0 <= v1 OR a@0 IN LIST [v2, v3, v4 ...] OR a@0 >= v6
WHEN 1: ... OR ... OR ...
WHEN 2 ...
WHEN 3 ...
For non-CASE dynamic filters, we may want to or the ranges / in lists:
ex. combine these two
DynamicFilterPhysicalExpr: a@0 >= 0 AND a@0 <= 5 OR IN LIST [10, 11]
DynamicFilterPhysicalExpr: a@0 >= 1 AND a@0 <= 10 OR IN LIST [12, 13]
into this
DynamicFilterPhysicalExpr: a@0 >= 0 AND a@0 <= 10 OR IN LIST [10, 11, 12, 13]
It would be nice if datafusion supported these "merge" operations natively rather than users having to perform "surgery" on PhysicalExpr. Doing it on PhysicalExpr is a bit brittle. If the PhysicalExprs used in dynamic filters were to change, then datafusion-distributed (and perhaps ballista etc.) would have to deal with that. The definition of "merge" actually depends on if the dynamic filter is partition aware (CASE expr vs no CASE) and the type of partitioning (range vs hash).
┌───────────────────────────┐ ┌───────────────────────────┐
│ ... │ │ ... │
┌───────┐ │ ┌───────────────────────┐│ │ ┌───────────────────────┐ │
│Stage 3│ │ │ NetworkShuffleExec ││ │ │ NetworkShuffleExec │ │
└───────┘ │ │ ││ │ │ │ │
│ └───────────────────────┘│ │ └───────────────────────┘ │
└───────────────────────────┘ └───────────────────────────┘
▲ ▲
│ │
┌──────────────────────────────────────────────────────┴────────────┬───────────────┴───────────────────────────────────────────────────┐
│ │ │
┌─────────────────────────────────┴────────────────────────────┐ ┌─────────────────────────────────┴────────────────────────────┐ ┌─────────────────────────────────┴───────────────────────────┐
│ ┌───────────────────────┐ │ │ ┌───────────────────────┐ │ │ ┌───────────────────────┐ │
│ │ RepartitionExec │ │ │ │ RepartitionExec │ │ │ │ RepartitionExec │ │
│ │ │ │ │ │ │ │ │ │ │ │
│ └───────────────────────┘ │ │ └───────────────────────┘ │ │ └───────────────────────┘ │
│ ┌────────────────────────┐ │ │ ┌────────────────────────┐ │ │ ┌────────────────────────┐ │
│ │ HashJoinExec: │ │ │ │ HashJoinExec: │ │ │ │ HashJoinExec: │ │
│ │ mode=DoesNotMatter │ │ │ │ mode=DoesNotMatter │ │ │ │ mode=DoesNotMatter │ │
│ └────────────────────────┘ │ │ └────────────────────────┘ │ │ └────────────────────────┘ │
┌───────┐ │ ▲ ▲ │ │ ▲ ▲ │ │ ▲ ▲ │
│Stage 2│ │ │ │ │ │ │ │ │ │ │ │ │
└───────┘ │ ┌──────────┘ └───────────┐ │ │ ┌──────────┘ └───────────┐ │ │ ┌──────────┘ └───────────┐ │
│ │ │ │ │ │ │ │ │ │ │ │
│ ┌──────────────────────────┐ ┌──────────────────────────┐ │ │ ┌──────────────────────────┐ ┌──────────────────────────┐ │ │ ┌──────────────────────────┐ ┌──────────────────────────┐ │
│ │ DataSourceExec │ │ NetworkShuffleExec │ │ │ │ DataSourceExec │ │ NetworkShuffleExec │ │ │ │ DataSourceExec │ │ NetworkShuffleExec │ │
│ │ │ │ │ │ │ │ │ │ │ │ │ │ │ │ │ │
│ └──────────────────────────┘ └──────────────────────────┘ │ │ └──────────────────────────┘ └──────────────────────────┘ │ │ └──────────────────────────┘ └──────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘ └──────────────────────────────────────────────────────────────┘ └─────────────────────────────────────────────────────────────┘
Task 1 Runs partitions [0,4) Task 2 Runs partitions [4,8) Task 3 Runs partitions [8,12)
▲ ▲ ▲
│ │ │
└──────────────────────────────────────────────┬───────────────────┴────────────────┬──────────────────────────────────────────────────┘
│ │
│ │
┌──────────────────────────────┐ ┌──────────────────────────────┐
│ ┌──────────────────────────┐ │ │ ┌──────────────────────────┐ │
│ │ RepartitionExec │ │ │ │ RepartitionExec │ │
│ │ partitioning=Hash(a, 12) │ │ │ │ partitioning=Hash(a, 12) │ │
│ └──────────────────────────┘ │ │ └──────────────────────────┘ │
│ ▲ │ │ ▲ │
│ │ │ │ │ │
┌───────┐ │ ... random operators │ │ ... random operators │
│Stage 1│ │ │ │ │
└───────┘ │ ▲ │ │ ▲ │
│ │ │ │ │ │
│ ┌─────────────┴────────────┐ │ │ ┌─────────────┴────────────┐ │
│ │ DataSourceExec: │ │ │ │ DataSourceExec: │ │
│ │ Partitioning=Unknown(8) │ │ │ │ Partitioning=Unknown(8) │ │
│ └──────────────────────────┘ │ │ └──────────────────────────┘ │
└──────────────────────────────┘ └──────────────────────────────┘
Task 1 Runs partitions [0, 4) Task 2 Runs partitions [4, 8)
Describe the solution you'd like
A very highlevel design
Today dynamic filters store one expression which is updated atomically:
struct Inner {
expression_id: u64,.
generation: u64,
// The actual dynamic filter expression
expr: Arc<dyn PhysicalExpr>,
is_complete: bool,
}
impl DynamicFilterPhysicalExpr {
// Atomically update the expression
pub fn update(&self, new_expr: Arc<dyn PhysicalExpr>) -> Result<()>;
// Atomically get the current expression
pub fn current(&self) -> Result<Arc<dyn PhysicalExpr>>;
}
It would be interesting to make them more partition-aware by:
- store one expression per partition when the dynamic filter is "partition aware", otherwise just store one
- change
update() to update the expression for a specific partition
- update
current() to return a generated CASE expression
- implement a
union() operation
struct Inner {
expression_id: u64,
generation: u64,
// Instead of one expr, we may have one expr per partition
expr: LiveFilterExpr,
lowered_expr: Arc<dyn PhysicalExpr>,
is_complete: bool,
}
// The actual dynamic filter expression
enum LiveFilterExpr {
Global(Arc<dyn PhysicalExpr>),
Partitioned(PartitionedFilterExpr),
}
impl LiveFilterExpr {
// Generates a new expression by ORing.
// For LiveFilterExpr::Global, we can OR the cases together.
// For LiveFilterExpr::Partitioned, we can modify the cases as needed (ex. we do a CASE-wise OR for hash partitioned filters).
// We may also alter the behavior depending on if the partitioning is hash vs range.
pub fn union(&self, other: LiveFilterExpr) -> Result<()>;
}
struct PartitionedFilterExpr {
partitioning: Partitioning, // ex. Partitioning=Hash(column_a, 12)
partition_expr: Arc<dyn PhysicalExpr>, // ex. hash(column_a)
cases: BTreeMap<u64, Arc<dyn PhysicalExpr>>, // map of partition id to filter expr
}
impl DynamicFilterPhysicalExpr {
// Update takes an optional partition id.
// Error if called with a partition id when the dynamic filter is not partition-aware and vice versa.
pub fn update(&self, expr: Arc<dyn PhysicalExpr>, partition_id: Option<u64>) -> Result<()>;
// Method is the same but now generates CASE expressions
pub fn current(&self) -> Result<Arc<dyn PhysicalExpr>>;
}
We would probably want to refactor datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs (~700 LoC) to migrate the partition specific logic into the DynamicFilterPhysicalExpr itself.
Describe alternatives you've considered
No response
Additional context
No response
Is your feature request related to a problem or challenge?
I think this would help solve two problems
1) Remove CASE hash(expr)
Related: #23376 (comment)
To support cases like below where the partitioning of the dynamic filter producer (hash join) is different than the consumer (data source), the hash join uses a
CASE hash(expr)to "bake" partitioning in to the join.What if the partitioning of the data source is the same as the hash join? In this case, the
CASE hash(expr)is pure overhead.It would be nice if the
DataSourceExeccould consume per-partition dynamic filters. I think mostDataSourceExechave to downcastPhysicalExprtoDynamicFilterPhysicalExprto consume them anyways. If theDynamicFilterPhysicalExprstored expressions per partition, theDataSourceExeccould apply per-partition filters rather than use theCASE hash(expr)2) Distributed Execution
More tails in the RFC: datafusion-contrib/datafusion-distributed#553
In datafusion-distributed, I'm interested in implementing dynamic filtering similarly to Trino.
HashJoinExec,SortExec,AggregateExec) will send dynamic filter updates to the coordinator.DataSourceExec).For correctness, the coordinator must wait for all dynamic filter updates from all workers before sending them to the consumers, otherwise, the consumers may prune rows incorrectly.
In the diagram below, each task in stage 2 creates a dynamic filter, so there's 3 in total:
Task 1:
Task 2:
Task 3:
The plan is to merge them together.
Option 1: OR them
Option 2: Case-wise OR
For non-CASE dynamic filters, we may want to or the ranges / in lists:
It would be nice if datafusion supported these "merge" operations natively rather than users having to perform "surgery" on
PhysicalExpr. Doing it onPhysicalExpris a bit brittle. If thePhysicalExprs used in dynamic filters were to change, then datafusion-distributed (and perhaps ballista etc.) would have to deal with that. The definition of "merge" actually depends on if the dynamic filter is partition aware (CASE expr vs no CASE) and the type of partitioning (range vs hash).Describe the solution you'd like
A very highlevel design
Today dynamic filters store one expression which is updated atomically:
It would be interesting to make them more partition-aware by:
update()to update the expression for a specific partitioncurrent()to return a generatedCASEexpressionunion()operationWe would probably want to refactor
datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs(~700 LoC) to migrate the partition specific logic into theDynamicFilterPhysicalExpritself.Describe alternatives you've considered
No response
Additional context
No response