weiqingy commented on code in PR #28973:
URL: https://github.com/apache/flink/pull/28973#discussion_r4232976569


##########
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:
   Good catch. `LogicalType.equals` compares nullability, so a leaf field 
changing nullable to NOT NULL is already rejected, but `bothRow` 
short-circuited that check for a ROW field and the nested recursion only ever 
compares the ROW's children, never the ROW's own nullability. So the narrowing 
was accepted and `getNewRowData` then left the slot null, putting a null into a 
field declared NOT NULL with nothing to catch it.
   
   Fixed in 46ff1d68a21 by settling nullability before delegating, which also 
closes the same hole on the positional path. Three tests added, and I checked 
they fail without the fix.
   



##########
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:
   Correct, that path had no coverage at all. The existing tests use a fake 
backend that overrides `getOrCreateKeyedState`, so the production arming line 
never ran.
   
   Added both capability cases against a real `AbstractKeyedStateBackend` in 
46ff1d68a21, using the existing mock backend with a new opt-in builder flag. 
Verified the armed case fails if the arming line is reverted.
   



##########
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:
   Agreed that the per-entry work is real: `migrate` calls 
`restoreSerializer()` on both snapshots for every restored value, and the 
position mapping is rebuilt per row and per nested row.
   
   I would rather not do it in this PR. Caching puts mutable state on a 
snapshot object whose thread-safety contract I have not pinned down, and the 
win is a constant factor on a path that is correct today. FLINK-40299 adds the 
end-to-end RocksDB restore coverage, which is where this can be measured rather 
than guessed at. I will pick it up there.
   



##########
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:
   I would rather not log it here. Effective activation is not a single 
boolean: it depends on the backend capability, the serializer shape, and the 
descriptor, so one line at keyed-state init cannot state it truthfully, and it 
would fire once per state per subtask.
   
   The underlying problem is real though. "I set the option and nothing 
happened" is the failure mode I most expect, including the case where RocksDB 
is configured but the changelog backend is also on. That is covered by the user 
documentation in FLINK-40299, which spells out when the option applies and when 
it silently does not.
   



-- 
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