Weiqing Yang created FLINK-40932:
------------------------------------

             Summary: Arm state schema evolution for ListState and MapState of 
RowData
                 Key: FLINK-40932
                 URL: https://issues.apache.org/jira/browse/FLINK-40932
             Project: Flink
          Issue Type: Sub-task
          Components: API / Type Serialization System, Runtime / State Backends
            Reporter: Weiqing Yang
            Assignee: Weiqing Yang


FLIP-527 covers list-element and map-value migration as well as value state. 
From the FLIP:

> list state (whose elements are stored individually) and map state (whose 
> values are stored per entry) are migrated by applying the same migrate method 
> to each element or value as it is read — there is no separate element hook.

FLINK-40298 arms only a state's own value serializer, so `ListState<RowData>` 
and `MapState<K, RowData>` currently fail closed: an evolving schema is 
rejected on restore rather than migrated. This sub-task extends the arming seam 
one structural level to cover them.

The runtime already descends to exactly these two serializers during migration 
(`RocksDBListState#migrateSerializedValue` unwraps the element serializer, 
`RocksDBMapState#migrateSerializedValue` unwraps the map value serializer), so 
this is an arming change rather than a new migration path.

Scope:

- Extend the arming helper with a `ListSerializer` branch (arm the element 
serializer) and a `MapSerializer` branch (arm the value serializer, never the 
key serializer).
- The descent is exactly one structural level. A `RowData` serializer two or 
more levels down, as in the interval-join shape `MapState<Long, 
List<Tuple2<RowData, Boolean>>>`, must remain unarmed and continue to fail 
closed.
- Map keys are never armed; key-schema changes remain incompatible.
- Fail closed at the consuming end. Extending the descent means an armed 
serializer can now reach operator state and broadcast state, because a 
`StateDescriptor` caches its serializer per instance and that cache is shared 
between the keyed and non-keyed stores. `DefaultOperatorStateBackend` accepts 
`compatibleAfterMigration` and never migrates, so it must reject an armed 
serializer instead. This adds a small change in `flink-runtime` alongside the 
`flink-core` one.
- Tests covering list-element and map-value migration, a test asserting the 
two-level shape is still rejected, and tests for the operator-state and 
broadcast-state rejection.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to