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]

Reply via email to