gene-bordegaray opened a new issue, #25483:
URL: https://github.com/apache/datafusion/issues/25483
Epic: #25421
`RangePartitioning` provides a separated, ordered key space using an
ordering and split points.
Currently, DF can determine that two range layouts are co-partitioned when
they have the same partition count, split points, and sort options ( see
`compatible_co_partitioning_layout` in
`datafusion/physical-plan/src/distribution_requirements.rs`).
Similarly, `InterleaveExec` requires exact partitioning equality. DF does
not describe
how two different but compatible layouts relate.
Lets look at a source and target layout to understand the relationships
between the two:
```text
Source layout Target layout
S0: (-inf, 10) ───┐
├─────────────────▶ T0: (-inf, 20)
S1: [10, 20) ─────┘
S2: [20, 30) ─────┐
├─────────────────▶ T1: [20, +inf)
S3: [30, +inf) ───┘
```
The two layouts define the space at different granularities: source
partitions define the space with 4 while the target defines it with 2. What is
important to not is that in order to transform the souce layout to the target,
it is not necessary for every source partition to have a relationship with
every target partition. Rather, every source partition has exactly one target
partition. Thus, rows do not need to be evaluated or redistributed individually
and adjacent source partitions can be grouped.
Nevertheless, the layouts are not equal, so the optimizer today cannot use
this relationship.
This missing information affects range repartitioning. `RepartitionExec`
constructs channels from every input partition to every output partition (an
all-to-all relationship). Given six inputs and three outputs this looks like:
```text
┌──────────────> C0
P0 ───────────────┼──────────────▶ C1
└──────────────> C2
┌──────────────> C0
P1 ───────────────┼──────────────▶ C1
└──────────────> C2
...
┌──────────────> C0
P5 ───────────────┼──────────────▶ C1
└──────────────> C2
6 inputs * 3 outputs = 18 relationships
```
When input range metadata can that only a subset of these edges can carry
rows the remaining channels still exist. Furthermore, an output will keep
channels open from inputs that can never produce rows for it and can't complete
until evey unrelated producer finishes.
## Desired behavior
DF should be able to identify the relationships between two range layouts
and return the corresponding partition mapping. This can be used to reduce
repartition overhead and make cleaner, more comprehensible optimizer decisions.
I propose the initial relationships should be:
1. **Equal**: Every source partition matches one target partition.
2. **SourceFiner**: The source is more granular of the same space. Adjacent
source partitions can be grouped to form target partitions.
3. **TargetFiner**: The target is more granular than the source. Only source
partitions crossed by a target range boundary require row splitting.
4. **Unproven**: Cannot prove a useful relationship.
I would prefer to start by being very strict here. Starting with this
focused set, proving and testing them thoroughly and if we cannot prove a
relationship fallback.
## Why this is good
This will help single node DF by eliminating lots of overhead when dealing
with range layouts by skipping repartitions and row level comparisions in many
scenarios.
This is only amplified in a distributed context for projects like
distributed datafusion and Ballista where minimixaing these relationships turns
into eliminating lots of expensive network coordination and transfers.
## Proposed solution
### 1. Add a range-layout relationship API
The API should return both classification and a mapping:
```rust
enum RangeLayoutRelationshipKind {
Equal,
SourceFiner,
TargetFiner,
Unproven,
}
struct RangeLayoutRelationship {
kind: RangeLayoutRelationshipKind,
source_to_targets: Vec<Range<>>,
target_to_sources: Vec<Range<>>,
}
**NOTE: range here is the range of partitions blogining to mapping at index
i not range partitioning**
```
The mapping can be computed via the two split-point lists.
### 2. Use the relationship during distribution enforcement
Distribution enforcement should use the relationship like this:
1. **Equal**: Already co-partitioned, reuse existing layout.
2. **SourceFiner**: Group adjacent source partitions (or for first PR could
just still insert repartitions on to limit scope).
3. **Target Finer**: Split across crossed partitions
4. **Unproven**: Not co-partitioned, fallback to repartition.
### 3. Make `RepartitionExec` use channel mapping
When a proven input-to-output mapping is available:
1. Create channels and spill state only for mapped `(input, output)` pairs
2. If an input only has one destination just send whole batches without
computing per-row partition indices.
3. If an input has multiple destinations use the range router.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]