Copilot commented on code in PR #28973:
URL: https://github.com/apache/flink/pull/28973#discussion_r4215971230
##########
flink-runtime/src/main/java/org/apache/flink/runtime/state/AbstractKeyedStateBackend.java:
##########
@@ -387,7 +390,7 @@ public <N, S extends State, V> S getOrCreateKeyedState(
InternalKvState<K, ?, ?> kvState =
keyValueStatesByName.get(stateDescriptor.getName());
if (kvState == null) {
if (!stateDescriptor.isSerializerInitialized()) {
- stateDescriptor.initializeSerializerUnlessSet(executionConfig);
+
stateDescriptor.initializeSerializerUnlessSet(stateValueSerializerFactory());
Review Comment:
The new arming path for descriptors registered directly through
`AbstractKeyedStateBackend#getOrCreateKeyedState` is not exercised by the added
arming tests; those tests cover `DefaultKeyedStateStore` and
`StreamingRuntimeContext` only. Table/window operators call this backend method
directly, so a regression here could leave those serializers unarmed while all
current tests still pass. Add a focused backend test covering both capability
values and asserting the descriptor's final serializer.
##########
flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/runtime/typeutils/RowDataSerializer.java:
##########
@@ -387,25 +483,224 @@ public TypeSerializerSchemaCompatibility<RowData>
resolveSchemaCompatibility(
RowDataSerializerSnapshot oldRowDataSerializerSnapshot =
(RowDataSerializerSnapshot) oldSerializerSnapshot;
- if (!Arrays.equals(types, oldRowDataSerializerSnapshot.types)) {
+ // A side that carries no names cannot disagree with anything:
without names a reorder
+ // is undetectable, so identical types keep meaning "compatible as
is" there, exactly as
+ // they do unarmed.
+ boolean namesDisagree =
+ stateSchemaEvolutionEnabled
+ && fieldNames != null
+ && oldRowDataSerializerSnapshot.fieldNames != null
+ && !Arrays.equals(fieldNames,
oldRowDataSerializerSnapshot.fieldNames);
+
+ // Identical positional layout: the nested composite path. Equal
types at equal
+ // positions say nothing about which field is which, so an armed
resolution has to see
+ // the names agree as well before it can treat the layout as
unchanged -- otherwise a
+ // reorder or a rename among same-typed fields resolves here as
needing no migration
+ // and leaves every value sitting under a neighbour's name.
+ if (Arrays.equals(types, oldRowDataSerializerSnapshot.types) &&
!namesDisagree) {
+
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData>
+ intermediateResult =
+ CompositeTypeSerializerUtil
+
.constructIntermediateCompatibilityResult(
+
nestedSerializersSnapshotDelegate
+
.getNestedSerializerSnapshots(),
+ oldRowDataSerializerSnapshot
+
.nestedSerializersSnapshotDelegate
+
.getNestedSerializerSnapshots());
+
+ if
(intermediateResult.isCompatibleWithReconfiguredSerializer()) {
+ RowDataSerializer reconfiguredCompositeSerializer =
restoreSerializer();
+ return
TypeSerializerSchemaCompatibility.compatibleWithReconfiguredSerializer(
+ reconfiguredCompositeSerializer);
+ }
+
+ return intermediateResult.getFinalResult();
+ }
+
+ if (!stateSchemaEvolutionEnabled) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+
+ // The new side must carry field names. A name-less new serializer
is reachable -- a
+ // structured type resolves to one -- and admitting it would open
a permanent name-less
+ // evolution channel rather than a ramp for savepoints taken
before names were stored.
+ if (fieldNames == null) {
return TypeSerializerSchemaCompatibility.incompatible();
}
-
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData>
- intermediateResult =
-
CompositeTypeSerializerUtil.constructIntermediateCompatibilityResult(
- nestedSerializersSnapshotDelegate
- .getNestedSerializerSnapshots(),
-
oldRowDataSerializerSnapshot.nestedSerializersSnapshotDelegate
- .getNestedSerializerSnapshots());
-
- if (intermediateResult.isCompatibleWithReconfiguredSerializer()) {
- RowDataSerializer reconfiguredCompositeSerializer =
restoreSerializer();
- return
TypeSerializerSchemaCompatibility.compatibleWithReconfiguredSerializer(
- reconfiguredCompositeSerializer);
+ return oldRowDataSerializerSnapshot.fieldNames != null
+ ? checkNameBasedEvolution(oldRowDataSerializerSnapshot)
+ : checkPositionalEvolution(oldRowDataSerializerSnapshot);
+ }
+
+ private TypeSerializerSchemaCompatibility<RowData>
checkNameBasedEvolution(
+ RowDataSerializerSnapshot oldSnapshot) {
+ int[] oldToNew = buildNameMapping(oldSnapshot.fieldNames,
this.fieldNames);
+ int[] newToOld = buildNameMapping(this.fieldNames,
oldSnapshot.fieldNames);
+
+ // (A) Every new-only field (no matching old field) must be
nullable.
+ for (int newPos = 0; newPos < newToOld.length; newPos++) {
+ if (newToOld[newPos] == -1 && !types[newPos].isNullable()) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
}
- return intermediateResult.getFinalResult();
+ // (B) Every old field must survive with a compatible type, and
nested snapshots are
+ // aligned old->new so nested ROW evolution can recurse. Leaf
(non-ROW) fields
+ // require an exactly equal type; ROW fields defer to the
nested recursion in (C).
+ TypeSerializerSnapshot<?>[] newNested =
+
nestedSerializersSnapshotDelegate.getNestedSerializerSnapshots();
+ TypeSerializerSnapshot<?>[] alignedNewNested =
+ new TypeSerializerSnapshot<?>[oldSnapshot.types.length];
+ for (int oldPos = 0; oldPos < oldToNew.length; oldPos++) {
+ int newPos = oldToNew[oldPos];
+ if (newPos == -1) {
+ return TypeSerializerSchemaCompatibility.incompatible();
// field removed
+ }
+ LogicalType oldType = oldSnapshot.types[oldPos];
+ LogicalType newType = types[newPos];
+ if (!bothRow(oldType, newType) && !oldType.equals(newType)) {
+ return TypeSerializerSchemaCompatibility.incompatible();
// leaf type changed
+ }
+ alignedNewNested[oldPos] = newNested[newPos];
+ }
+
+ // (C) Recurse into the aligned nested snapshot pairs.
+ return resolveAlignedNested(alignedNewNested, oldSnapshot);
+ }
+
+ /**
+ * Resolves against a prior snapshot that carries no field names,
matching fields by
+ * position.
+ *
+ * <p>Position is a stable identity only for an append. An insertion
in the middle is
+ * indistinguishable from a retype plus an append, and the two demand
opposite migrations,
+ * so the old layout has to be a prefix of the new one.
+ */
+ private TypeSerializerSchemaCompatibility<RowData>
checkPositionalEvolution(
+ RowDataSerializerSnapshot oldSnapshot) {
+ if (types.length < oldSnapshot.types.length) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ for (int i = 0; i < oldSnapshot.types.length; i++) {
+ LogicalType oldType = oldSnapshot.types[i];
+ LogicalType newType = types[i];
+ if (!bothRow(oldType, newType) && !oldType.equals(newType)) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ }
+ for (int i = oldSnapshot.types.length; i < types.length; i++) {
+ if (!types[i].isNullable()) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ }
+
+ // constructIntermediateCompatibilityResult requires both arrays
to have the same
+ // length, so only the prefix the old layout covers is handed to
it.
+ TypeSerializerSnapshot<?>[] alignedNewNested =
+ Arrays.copyOf(
+
nestedSerializersSnapshotDelegate.getNestedSerializerSnapshots(),
+ oldSnapshot.types.length);
+ return resolveAlignedNested(alignedNewNested, oldSnapshot);
+ }
+
+ private TypeSerializerSchemaCompatibility<RowData>
resolveAlignedNested(
+ TypeSerializerSnapshot<?>[] alignedNewNested,
+ RowDataSerializerSnapshot oldSnapshot) {
+
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData> nested =
+
CompositeTypeSerializerUtil.constructIntermediateCompatibilityResult(
+ alignedNewNested,
+ oldSnapshot.nestedSerializersSnapshotDelegate
+ .getNestedSerializerSnapshots());
+ // A reconfigured nested serializer is deliberately not propagated
here, unlike on the
+ // identical-layout path. Reconfiguration exists so a new
serializer can read old bytes;
+ // once the values have been remapped there are no old bytes left,
because the migrated
+ // row is re-encoded by the state's own new serializer.
+ return nested.isIncompatible()
+ ? TypeSerializerSchemaCompatibility.incompatible()
+ :
TypeSerializerSchemaCompatibility.compatibleAfterMigration();
+ }
+
+ private static boolean bothRow(LogicalType oldType, LogicalType
newType) {
+ return oldType.getTypeRoot() == LogicalTypeRoot.ROW
+ && newType.getTypeRoot() == LogicalTypeRoot.ROW;
+ }
Review Comment:
`bothRow` skips full type equality for every pair of ROWs, including when
the field itself changes from nullable to `NOT NULL`. That change is therefore
accepted, but `getNewRowData` preserves an existing null nested value and
reserializes it into the new non-null schema. Require the ROW container
nullability to match before delegating its children to recursive compatibility
(consistent with the exact-equality check used for leaf fields).
##########
flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/runtime/typeutils/RowDataSerializer.java:
##########
@@ -387,25 +483,224 @@ public TypeSerializerSchemaCompatibility<RowData>
resolveSchemaCompatibility(
RowDataSerializerSnapshot oldRowDataSerializerSnapshot =
(RowDataSerializerSnapshot) oldSerializerSnapshot;
- if (!Arrays.equals(types, oldRowDataSerializerSnapshot.types)) {
+ // A side that carries no names cannot disagree with anything:
without names a reorder
+ // is undetectable, so identical types keep meaning "compatible as
is" there, exactly as
+ // they do unarmed.
+ boolean namesDisagree =
+ stateSchemaEvolutionEnabled
+ && fieldNames != null
+ && oldRowDataSerializerSnapshot.fieldNames != null
+ && !Arrays.equals(fieldNames,
oldRowDataSerializerSnapshot.fieldNames);
+
+ // Identical positional layout: the nested composite path. Equal
types at equal
+ // positions say nothing about which field is which, so an armed
resolution has to see
+ // the names agree as well before it can treat the layout as
unchanged -- otherwise a
+ // reorder or a rename among same-typed fields resolves here as
needing no migration
+ // and leaves every value sitting under a neighbour's name.
+ if (Arrays.equals(types, oldRowDataSerializerSnapshot.types) &&
!namesDisagree) {
+
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData>
+ intermediateResult =
+ CompositeTypeSerializerUtil
+
.constructIntermediateCompatibilityResult(
+
nestedSerializersSnapshotDelegate
+
.getNestedSerializerSnapshots(),
+ oldRowDataSerializerSnapshot
+
.nestedSerializersSnapshotDelegate
+
.getNestedSerializerSnapshots());
+
+ if
(intermediateResult.isCompatibleWithReconfiguredSerializer()) {
+ RowDataSerializer reconfiguredCompositeSerializer =
restoreSerializer();
+ return
TypeSerializerSchemaCompatibility.compatibleWithReconfiguredSerializer(
+ reconfiguredCompositeSerializer);
+ }
+
+ return intermediateResult.getFinalResult();
+ }
+
+ if (!stateSchemaEvolutionEnabled) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+
+ // The new side must carry field names. A name-less new serializer
is reachable -- a
+ // structured type resolves to one -- and admitting it would open
a permanent name-less
+ // evolution channel rather than a ramp for savepoints taken
before names were stored.
+ if (fieldNames == null) {
return TypeSerializerSchemaCompatibility.incompatible();
}
-
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData>
- intermediateResult =
-
CompositeTypeSerializerUtil.constructIntermediateCompatibilityResult(
- nestedSerializersSnapshotDelegate
- .getNestedSerializerSnapshots(),
-
oldRowDataSerializerSnapshot.nestedSerializersSnapshotDelegate
- .getNestedSerializerSnapshots());
-
- if (intermediateResult.isCompatibleWithReconfiguredSerializer()) {
- RowDataSerializer reconfiguredCompositeSerializer =
restoreSerializer();
- return
TypeSerializerSchemaCompatibility.compatibleWithReconfiguredSerializer(
- reconfiguredCompositeSerializer);
+ return oldRowDataSerializerSnapshot.fieldNames != null
+ ? checkNameBasedEvolution(oldRowDataSerializerSnapshot)
+ : checkPositionalEvolution(oldRowDataSerializerSnapshot);
+ }
+
+ private TypeSerializerSchemaCompatibility<RowData>
checkNameBasedEvolution(
+ RowDataSerializerSnapshot oldSnapshot) {
+ int[] oldToNew = buildNameMapping(oldSnapshot.fieldNames,
this.fieldNames);
+ int[] newToOld = buildNameMapping(this.fieldNames,
oldSnapshot.fieldNames);
+
+ // (A) Every new-only field (no matching old field) must be
nullable.
+ for (int newPos = 0; newPos < newToOld.length; newPos++) {
+ if (newToOld[newPos] == -1 && !types[newPos].isNullable()) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
}
- return intermediateResult.getFinalResult();
+ // (B) Every old field must survive with a compatible type, and
nested snapshots are
+ // aligned old->new so nested ROW evolution can recurse. Leaf
(non-ROW) fields
+ // require an exactly equal type; ROW fields defer to the
nested recursion in (C).
+ TypeSerializerSnapshot<?>[] newNested =
+
nestedSerializersSnapshotDelegate.getNestedSerializerSnapshots();
+ TypeSerializerSnapshot<?>[] alignedNewNested =
+ new TypeSerializerSnapshot<?>[oldSnapshot.types.length];
+ for (int oldPos = 0; oldPos < oldToNew.length; oldPos++) {
+ int newPos = oldToNew[oldPos];
+ if (newPos == -1) {
+ return TypeSerializerSchemaCompatibility.incompatible();
// field removed
+ }
+ LogicalType oldType = oldSnapshot.types[oldPos];
+ LogicalType newType = types[newPos];
+ if (!bothRow(oldType, newType) && !oldType.equals(newType)) {
+ return TypeSerializerSchemaCompatibility.incompatible();
// leaf type changed
+ }
+ alignedNewNested[oldPos] = newNested[newPos];
+ }
+
+ // (C) Recurse into the aligned nested snapshot pairs.
+ return resolveAlignedNested(alignedNewNested, oldSnapshot);
+ }
+
+ /**
+ * Resolves against a prior snapshot that carries no field names,
matching fields by
+ * position.
+ *
+ * <p>Position is a stable identity only for an append. An insertion
in the middle is
+ * indistinguishable from a retype plus an append, and the two demand
opposite migrations,
+ * so the old layout has to be a prefix of the new one.
+ */
+ private TypeSerializerSchemaCompatibility<RowData>
checkPositionalEvolution(
+ RowDataSerializerSnapshot oldSnapshot) {
+ if (types.length < oldSnapshot.types.length) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ for (int i = 0; i < oldSnapshot.types.length; i++) {
+ LogicalType oldType = oldSnapshot.types[i];
+ LogicalType newType = types[i];
+ if (!bothRow(oldType, newType) && !oldType.equals(newType)) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ }
+ for (int i = oldSnapshot.types.length; i < types.length; i++) {
+ if (!types[i].isNullable()) {
+ return TypeSerializerSchemaCompatibility.incompatible();
+ }
+ }
+
+ // constructIntermediateCompatibilityResult requires both arrays
to have the same
+ // length, so only the prefix the old layout covers is handed to
it.
+ TypeSerializerSnapshot<?>[] alignedNewNested =
+ Arrays.copyOf(
+
nestedSerializersSnapshotDelegate.getNestedSerializerSnapshots(),
+ oldSnapshot.types.length);
+ return resolveAlignedNested(alignedNewNested, oldSnapshot);
+ }
+
+ private TypeSerializerSchemaCompatibility<RowData>
resolveAlignedNested(
+ TypeSerializerSnapshot<?>[] alignedNewNested,
+ RowDataSerializerSnapshot oldSnapshot) {
+
CompositeTypeSerializerUtil.IntermediateCompatibilityResult<RowData> nested =
+
CompositeTypeSerializerUtil.constructIntermediateCompatibilityResult(
+ alignedNewNested,
+ oldSnapshot.nestedSerializersSnapshotDelegate
+ .getNestedSerializerSnapshots());
+ // A reconfigured nested serializer is deliberately not propagated
here, unlike on the
+ // identical-layout path. Reconfiguration exists so a new
serializer can read old bytes;
+ // once the values have been remapped there are no old bytes left,
because the migrated
+ // row is re-encoded by the state's own new serializer.
+ return nested.isIncompatible()
+ ? TypeSerializerSchemaCompatibility.incompatible()
+ :
TypeSerializerSchemaCompatibility.compatibleAfterMigration();
+ }
+
+ private static boolean bothRow(LogicalType oldType, LogicalType
newType) {
+ return oldType.getTypeRoot() == LogicalTypeRoot.ROW
+ && newType.getTypeRoot() == LogicalTypeRoot.ROW;
+ }
+
+ @Override
+ public RowData migrate(
+ TypeSerializerSnapshot<RowData> oldSerializerSnapshot, RowData
value) {
+ if (value == null) {
+ return null;
+ }
+ // Runs once per restored entry, so a mismatch would otherwise
surface as a bare
+ // ClassCastException from deep inside a migration loop.
+ Preconditions.checkArgument(
+ oldSerializerSnapshot instanceof RowDataSerializerSnapshot,
+ "Cannot migrate RowData state from %s.",
+ oldSerializerSnapshot.getClass().getName());
+ RowDataSerializerSnapshot oldSnapshot =
+ (RowDataSerializerSnapshot) oldSerializerSnapshot;
+ return getNewRowData(value, oldSnapshot.restoreSerializer(),
restoreSerializer());
Review Comment:
This method runs once for every restored state entry, but each call
reconstructs both complete serializer trees; `getNewRowData` then rebuilds the
same name-to-position maps for every row and nested row. RocksDB invokes this
inside its full-state migration loop, so large states incur avoidable per-entry
schema-sized allocations and GC during recovery. Cache a migration plan
(restored getters/serializers, position mappings, and nested plans) per old/new
snapshot pair and apply only that plan per value.
##########
flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java:
##########
@@ -45,6 +45,23 @@ public class ExecutionConfigOptions {
// State Options
// ------------------------------------------------------------------------
+ @Documentation.TableOption(execMode = Documentation.ExecMode.STREAMING)
+ public static final ConfigOption<Boolean>
TABLE_EXEC_STATE_SCHEMA_EVOLUTION_ENABLED =
+ key("table.exec.state.schema-evolution.enabled")
+ .booleanType()
+ .defaultValue(false)
Review Comment:
This new feature flag has no INFO-level initialization log for its resolved
and effective state. Because enabling the option is further gated by backend
capability and serializer shape, a restore rejection currently gives operators
no startup evidence showing whether schema evolution was enabled, disabled, or
blocked by the backend. Log the option value together with the backend
capability/effective activation once when the keyed-state path initializes.
--
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]