java-topology/defects/flink/patch/flink-0006-multiple-input-node-members-set.patch

28 lines
1.7 KiB
Diff

# UNDF: UNDF-2026-000000761
# UNDF: (leave blank)
--- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/processor/MultipleInputNodeCreationProcessor.java
+++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/processor/MultipleInputNodeCreationProcessor.java
@@ -534,7 +534,7 @@ public class MultipleInputNodeCreationProcessor implements DAGProcessor {
// calculate the inputs of the multiple input node
List<Tuple3<ExecNode<?>, InputProperty, ExecEdge>> inputs = new ArrayList<>();
+ Set<ExecNodeWrapper> membersSet = new HashSet<>(group.members);
for (ExecNodeWrapper member : group.members) {
for (int i = 0; i < member.inputs.size(); i++) {
ExecNodeWrapper memberInput = member.inputs.get(i);
- if (group.members.contains(memberInput)) {
+ if (membersSet.contains(memberInput)) {
continue;
}
@@ -705,10 +705,12 @@ public class MultipleInputNodeCreationProcessor implements DAGProcessor {
Preconditions.checkNotNull(
root, "Multiple input group does not have a root. This is a bug.");
- Set<ExecNodeWrapper> sameGroupInputWrappers = new HashSet<>();
+ Set<ExecNodeWrapper> membersSet = new HashSet<>(members);
+ Set<ExecNodeWrapper> sameGroupInputWrappers = new HashSet<>();
for (ExecNodeWrapper inputWrapper : root.inputs) {
- if (members.contains(inputWrapper)) {
+ if (membersSet.contains(inputWrapper)) {
sameGroupInputWrappers.add(inputWrapper);
}
}