Skip to content

Commit db27926

Browse files
committed
Preserve window/ranking sort boundary under ready-first Collect
When window/row-number/ranking nodes are pushed down per region, propagate the GroupNode ordering scheme to each split child so the upper merge uses MergeSort instead of an unordered Collect. Otherwise ready-first Collect scrambles partition continuity and corrupts ROW_NUMBER/RANK/DENSE_RANK results.
1 parent 7fcaba3 commit db27926

1 file changed

Lines changed: 24 additions & 3 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java‎

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3385,7 +3385,14 @@ public List<PlanNode> visitWindowFunction(WindowNode node, PlanContext context)
33853385
node.setChild(collectNode);
33863386
return Collections.singletonList(node);
33873387
} else {
3388-
return splitForEachChild(node, childrenNodes);
3388+
// ready-first Collect 无序合并会打散窗口所需的 partition 连续,下推后必须用 MergeSort
3389+
// 按 GroupNode 的排序键(partition key + order key)全局归并。
3390+
List<PlanNode> result = splitForEachChild(node, childrenNodes);
3391+
OrderingScheme windowOrdering = ((GroupNode) node.getChild()).getOrderingScheme();
3392+
for (PlanNode child : result) {
3393+
nodeOrderingMap.put(child.getPlanNodeId(), windowOrdering);
3394+
}
3395+
return result;
33893396
}
33903397
}
33913398

@@ -3415,7 +3422,14 @@ public List<PlanNode> visitRowNumber(RowNumberNode node, PlanContext context) {
34153422
node.setChild(collectNode);
34163423
return Collections.singletonList(node);
34173424
} else {
3418-
return splitForEachChild(node, childrenNodes);
3425+
// RowNumber 编号依赖 partition 连续,下推后必须用 MergeSort 按 GroupNode 的排序键
3426+
// 全局归并,否则 ready-first Collect 无序合并会打散 partition,编号出错。
3427+
List<PlanNode> result = splitForEachChild(node, childrenNodes);
3428+
OrderingScheme rowNumberOrdering = ((GroupNode) node.getChild()).getOrderingScheme();
3429+
for (PlanNode child : result) {
3430+
nodeOrderingMap.put(child.getPlanNodeId(), rowNumberOrdering);
3431+
}
3432+
return result;
34193433
}
34203434
}
34213435

@@ -3454,7 +3468,14 @@ public List<PlanNode> visitTopKRanking(TopKRankingNode node, PlanContext context
34543468
node.setChild(collectNode);
34553469
return Collections.singletonList(node);
34563470
} else {
3457-
return splitForEachChild(node, childrenNodes);
3471+
// 排名结果依赖 ORDER BY 顺序,下推后必须用 MergeSort 按 GroupNode 的排序键全局归并,
3472+
// 否则 ready-first Collect 无序合并会打散 partition 连续,编号/排名出错。
3473+
List<PlanNode> result = splitForEachChild(node, childrenNodes);
3474+
OrderingScheme rankingOrdering = ((GroupNode) node.getChild()).getOrderingScheme();
3475+
for (PlanNode child : result) {
3476+
nodeOrderingMap.put(child.getPlanNodeId(), rankingOrdering);
3477+
}
3478+
return result;
34583479
}
34593480
}
34603481

0 commit comments

Comments
 (0)