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]