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)