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]