You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Copy file name to clipboardExpand all lines: datafusion/sqllogictest/test_files/range_partitioning.slt
+4-177Lines changed: 4 additions & 177 deletions
Original file line number
Diff line number
Diff line change
@@ -1051,13 +1051,11 @@ SELECT range_key, SUM(value) OVER (ORDER BY value) FROM range_partitioned ORDER
1051
1051
30 1060
1052
1052
35 1410
1053
1053
1054
-
statement ok
1055
-
reset datafusion.explain.physical_plan_only;
1054
+
1056
1055
1057
1056
##########
1058
1057
# TEST 30: PartitionedTopK on Range Partition Column
1059
-
# With subset threshold met and preserve-file disabled, Range([range_key])
1060
-
# satisfies the TopK partition key and avoids repartitioning.
1058
+
# Exact Range([range_key]) satisfies the TopK partition key and avoids repartitioning.
1061
1059
##########
1062
1060
1063
1061
statement ok
@@ -1075,11 +1073,6 @@ EXPLAIN SELECT * FROM (
1075
1073
FROM range_partitioned
1076
1074
) WHERE rn <= 1;
1077
1075
----
1078
-
logical_plan
1079
-
01)Projection: range_partitioned.range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1080
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1081
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
1085
1078
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
@@ -1105,8 +1098,8 @@ ORDER BY range_key;
1105
1098
1106
1099
##########
1107
1100
# TEST 31: PartitionedTopK on Non-Range Column
1108
-
# With subset threshold met and preserve-file disabled, partitioning on a non-range
1109
-
# key cannot reuse Range([range_key]) and requires hash repartitioning.
1101
+
# Partitioning on a non-range key cannot reuse Range([range_key]) and
1102
+
# requires hash repartitioning.
1110
1103
##########
1111
1104
1112
1105
statement ok
@@ -1121,11 +1114,6 @@ EXPLAIN SELECT * FROM (
1121
1114
FROM range_partitioned
1122
1115
) WHERE rn <= 1;
1123
1116
----
1124
-
logical_plan
1125
-
01)Projection: range_partitioned.non_range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1126
-
02)--Filter: row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1127
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[non_range_key@0 as non_range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
1131
1119
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
@@ -1162,11 +1150,6 @@ EXPLAIN SELECT * FROM (
1162
1150
FROM range_partitioned
1163
1151
) WHERE rn <= 1;
1164
1152
----
1165
-
logical_plan
1166
-
01)Projection: range_partitioned.range_key, range_partitioned.non_range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1167
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1168
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, non_range_key@1 as non_range_key, value@2 as value, row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as rn]
1172
1155
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
@@ -1190,40 +1173,6 @@ ORDER BY range_key, non_range_key;
1190
1173
35 2 350 1
1191
1174
1192
1175
1193
-
##########
1194
-
# TEST 33: Exact Range PartitionedTopK Below Subset Threshold
1195
-
# Even when subset satisfaction is disabled, exact Range([range_key])
1196
-
# satisfies PARTITION BY range_key when repartitioning would not increase
1197
-
# partition count.
1198
-
##########
1199
-
1200
-
statement ok
1201
-
set datafusion.execution.target_partitions = 4;
1202
-
1203
-
statement ok
1204
-
set datafusion.optimizer.subset_repartition_threshold = 5;
1205
-
1206
-
statement ok
1207
-
set datafusion.optimizer.preserve_file_partitions = 0;
1208
-
1209
-
query TT
1210
-
EXPLAIN SELECT * FROM (
1211
-
SELECT range_key, value, ROW_NUMBER() OVER (PARTITION BY range_key ORDER BY value DESC) as rn
1212
-
FROM range_partitioned
1213
-
) WHERE rn <= 1;
1214
-
----
1215
-
logical_plan
1216
-
01)Projection: range_partitioned.range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1217
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1218
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
1222
-
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
# TEST 34: Range Subset PartitionedTopK Rehashes Below Subset Threshold
1229
1178
# Range([range_key]) is only a subset of PARTITION BY (range_key, non_range_key),
@@ -1246,138 +1195,16 @@ EXPLAIN SELECT * FROM (
1246
1195
FROM range_partitioned
1247
1196
) WHERE rn <= 1;
1248
1197
----
1249
-
logical_plan
1250
-
01)Projection: range_partitioned.range_key, range_partitioned.non_range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1251
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1252
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, non_range_key@1 as non_range_key, value@2 as value, row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as rn]
1256
1200
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key, range_partitioned.non_range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
# TEST 35: PartitionedTopK Rehashes Below Subset Threshold
1264
-
# With subset threshold 5 and only 4 input partitions, planning repartitions
1265
-
# to increase parallelism instead of reusing Range partitioning.
1266
-
##########
1267
-
1268
-
statement ok
1269
-
set datafusion.execution.target_partitions = 5;
1270
-
1271
-
statement ok
1272
-
set datafusion.optimizer.subset_repartition_threshold = 5;
1273
-
1274
-
statement ok
1275
-
set datafusion.optimizer.preserve_file_partitions = 0;
1276
-
1277
-
query TT
1278
-
EXPLAIN SELECT * FROM (
1279
-
SELECT range_key, value, ROW_NUMBER() OVER (PARTITION BY range_key ORDER BY value DESC) as rn
1280
-
FROM range_partitioned
1281
-
) WHERE rn <= 1;
1282
-
----
1283
-
logical_plan
1284
-
01)Projection: range_partitioned.range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1285
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1286
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
1290
-
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
# TEST 36: PartitionedTopK Preserves Range When Preserve File Threshold Met
1304
-
# With preserve-file threshold 1 and 4 input partitions, Range is preserved
1305
-
# even though target_partitions is 5.
1306
-
##########
1307
-
1308
-
statement ok
1309
-
set datafusion.execution.target_partitions = 5;
1310
-
1311
-
statement ok
1312
-
set datafusion.optimizer.subset_repartition_threshold = 4;
1313
-
1314
-
statement ok
1315
-
set datafusion.optimizer.preserve_file_partitions = 1;
1316
-
1317
-
query TT
1318
-
EXPLAIN SELECT * FROM (
1319
-
SELECT range_key, value, ROW_NUMBER() OVER (PARTITION BY range_key ORDER BY value DESC) as rn
1320
-
FROM range_partitioned
1321
-
) WHERE rn <= 1;
1322
-
----
1323
-
logical_plan
1324
-
01)Projection: range_partitioned.range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1325
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1326
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
03)----BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
# TEST 37: PartitionedTopK Rehashes When Preserve File Threshold Not Met
1344
-
# With preserve-file threshold 5 and only 4 input partitions, planning can
1345
-
# repartition to increase parallelism.
1346
-
##########
1347
-
1348
-
statement ok
1349
-
set datafusion.execution.target_partitions = 5;
1350
-
1351
-
statement ok
1352
-
set datafusion.optimizer.subset_repartition_threshold = 4;
1353
-
1354
-
statement ok
1355
-
set datafusion.optimizer.preserve_file_partitions = 5;
1356
-
1357
-
query TT
1358
-
EXPLAIN SELECT * FROM (
1359
-
SELECT range_key, value, ROW_NUMBER() OVER (PARTITION BY range_key ORDER BY value DESC) as rn
1360
-
FROM range_partitioned
1361
-
) WHERE rn <= 1;
1362
-
----
1363
-
logical_plan
1364
-
01)Projection: range_partitioned.range_key, range_partitioned.value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn
1365
-
02)--Filter: row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW <= UInt64(1)
1366
-
03)----WindowAggr: windowExpr=[[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
01)ProjectionExec: expr=[range_key@0 as range_key, value@1 as value, row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as rn]
03)----BoundedWindowAggExec: wdw=[row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [range_partitioned.range_key] ORDER BY [range_partitioned.value DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
0 commit comments