28 lines
1.7 KiB
Diff
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);
|
|
}
|
|
}
|