twalthr commented on code in PR #29339:
URL: https://github.com/apache/flink/pull/29339#discussion_r4153077041


##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecJoin.java:
##########
@@ -287,4 +319,12 @@ protected Transformation<RowData> translateToPlanInternal(
         transform.setStateKeyType(leftSelect.getProducedType());
         return transform;
     }
+
+    /** Compiled plans of previous versions do not contain the changelog modes 
of the inputs. */
+    private static @Nullable RuntimeChangelogMode toRuntimeChangelogMode(

Review Comment:
   `RuntimeChangelogMode` comes from PTFs. We should only store minimal 
information in CompiledPlan. As far as I see it, only the "isUpsert" 
information is required? Let's only add this boolean flag. Just to double 
check: leftUpsertKeys is also set in retract?



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecJoin.java:
##########
@@ -100,6 +104,16 @@ public class StreamExecJoin extends ExecNodeBase<RowData>
     @JsonInclude(JsonInclude.Include.NON_DEFAULT)
     private final List<int[]> rightUpsertKeys;
 
+    @Nullable
+    @JsonProperty(FIELD_NAME_LEFT_INPUT_CHANGELOG_MODE)
+    @JsonInclude(JsonInclude.Include.NON_NULL)
+    private final ChangelogMode leftInputChangelogMode;

Review Comment:
   Since this is a CompiledPlan change, make sure that one RestoreTest covers 
the old plan representation and one the new one.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to