1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//! [`EnsureRequirements`] optimizer rule that enforces distribution and
//! sorting requirements together so that the two never invalidate each other.
//!
//! This rule replaces the separate `EnforceDistribution` + `EnforceSorting`
//! rules with a unified approach inspired by Apache Spark's `EnsureRequirements`
//! and Presto/Trino's `AddExchanges`.
//!
//! # Motivation
//!
//! The previous two-rule design (`EnforceDistribution` then `EnforceSorting`)
//! suffers from non-idempotent composition: `EnforceSorting`'s `pushdown_sorts`
//! can break distribution invariants established by `EnforceDistribution`,
//! because `SortExec.preserve_partitioning` couples sorting and distribution
//! decisions. See <https://github.com/apache/datafusion/issues/21973> for details.
//!
//! # Architecture
//!
//! `optimize` runs several tree traversals. The defining property of this
//! rule is **Phase 2**: a single combined bottom-up pass that resolves
//! distribution *and* sorting for each node together. The surrounding phases
//! are independent traversals (top-down join-key reorder, then several
//! follow-up sort/order rewrites). Some of those could be consolidated
//! further in a follow-up.
//!
//! ```text
//! EnsureRequirements::optimize(plan)
//! │
//! ├─ Phase 1: top-down join-key reorder (adjust_input_keys_ordering)
//! │
//! ├─ Phase 2: combined distribution + sorting (single bottom-up pass)
//! │ └─ For each node (bottom-up), for each child:
//! │ Step 1: ensure distribution requirement
//! │ └─ insert RepartitionExec / CoalescePartitionsExec /
//! │ SortPreservingMergeExec as needed
//! │ Step 2: ensure ordering requirement (distribution-aware)
//! │ └─ insert SortExec with the correct `preserve_partitioning`,
//! │ with SortPreservingMergeExec on top if needed
//! │
//! └─ Phase 3: small follow-up passes (bottom-up unless noted)
//! ├─ parallelize_sorts
//! ├─ replace_with_order_preserving_variants
//! ├─ pushdown_sorts (recursive walk)
//! └─ replace_with_partial_sort
//! ```
//!
//! # Key Properties
//!
//! - **Idempotent across the whole rule**: Running `EnsureRequirements`
//! twice produces the same plan. This is the property that fixes
//! <https://github.com/apache/datafusion/issues/21973>, where the old
//! two-rule pipeline could regress a parallel sort plan into a serial one
//! on pass 2.
//! - **Distribution before sorting**: For each child, distribution is
//! resolved before ordering, so sorting decisions always have full
//! distribution context.
//! - **Sort pushdown is implicit**: Phase 2 only adds `SortExec` where the
//! child doesn't already satisfy the ordering requirement, so sorts land
//! at the deepest valid position without a separate destructive pass.
//!
//! # Behavior: parallelism via repartitioning
//!
//! Phase 2 Step 1 inserts `RepartitionExec` to satisfy distribution
//! requirements. When configuration allows, it also increases parallelism by
//! repartitioning over otherwise-serial inputs. For example, given two
//! 1-partition inputs feeding an operator that can run with more
//! parallelism:
//!
//! ```text
//! ┌─────────────────────────────────┐
//! │ ExecutionPlan │
//! └─────────────────────────────────┘
//! ▲ ▲
//! │ │
//! ┌───────────┐ ┌───────────┐
//! │ batch A │ │ batch B │ Input: 2 partitions
//! └───────────┘ └───────────┘
//! ```
//!
//! `EnsureRequirements` inserts a `RepartitionExec` so the operator runs
//! with three partitions:
//!
//! ```text
//! ┌─────────────────────────────────┐
//! │ ExecutionPlan │ Input now has 3 partitions
//! └─────────────────────────────────┘
//! ▲ ▲ ▲
//! └──────┼───────┘
//! │
//! ┌─────────────────────────────────┐
//! │ RepartitionExec(3) │ batches are repartitioned
//! │ RoundRobin │
//! └─────────────────────────────────┘
//! ▲ ▲
//! ┌───────────┐ ┌───────────┐
//! │ batch A │ │ batch B │
//! └───────────┘ └───────────┘
//! ```
//!
//! # Behavior: joint distribution + sorting
//!
//! Resolving distribution and sorting together lets Phase 2 produce a
//! parallel sort plan in cases where the two-rule pipeline historically
//! risked a serial one. Given `Sort(DESC) ← Coalesce ← MultiPartitionSource`,
//! `EnsureRequirements` rewrites it into:
//!
//! ```text
//! SortPreservingMergeExec: [a DESC] (cheap k-way merge of sorted streams)
//! SortExec: [a DESC], preserve_partitioning=true (N sorts run in parallel)
//! MultiPartitionSource
//! ```
//!
//! Each input partition is sorted in parallel, then a `SortPreservingMergeExec`
//! at the top performs a cheap merge of pre-sorted streams. For TopK queries
//! (`fetch=K`), each parallel sort only keeps K rows per partition, so total
//! memory is `N × K` rather than coalescing the entire stream first.
//!
//! # Behavior: strictest distribution match for joins
//!
//! Distribution requirements are met in the strictest way. For example, a
//! hash join with keys `(a, b, c)` requires `Distribution(a, b, c)`. This
//! can in principle be satisfied by partitioning on any superset of any
//! subset of `(a, b, c)`, but this rule always partitions on the exact key
//! tuple `(a, b, c)`. This is sometimes more aggressive than strictly
//! necessary, but the strictest match helps avoid data skew in joins.
// Internal implementation modules. Re-exported from `crate` root for tests
// in `core/tests/physical_optimizer/{enforce_distribution,enforce_sorting}.rs`.
use Arc;
use cratePhysicalOptimizerRule;
use Result;
use ConfigOptions;
use ;
use ExecutionPlan;
/// Optimizer rule that enforces both distribution and sorting requirements.
///
/// This rule combines the functionality of `EnforceDistribution` and
/// `EnforceSorting` into a coordinated sequence where distribution is
/// always settled before sorting for each operator, preventing the
/// non-idempotent interactions between the two separate rules.
///
/// See [module level documentation](self) for more details.
// See tests in datafusion/core/tests/physical_optimizer/ensure_requirements.rs