weiqingy opened a new pull request, #29449:
URL: https://github.com/apache/flink/pull/29449

   This is the PR-3b step of the 
[FLIP-527](https://cwiki.apache.org/confluence/spaces/FLINK/pages/353601981/FLIP-527+State+Schema+Evolution+for+RowData)
 implementation, split into a stack of small, independently reviewable PRs 
under the umbrella issue 
[FLINK-37732](https://issues.apache.org/jira/browse/FLINK-37732). Landing order:
   
   | Step | Sub-task | Scope |
   |---|---|---|
   | PR-1 | [FLINK-40296](https://issues.apache.org/jira/browse/FLINK-40296) | 
Object-level `migrate` hook on `TypeSerializerSnapshot` (merged) |
   | PR-2 | [FLINK-40297](https://issues.apache.org/jira/browse/FLINK-40297) | 
Route TTL-aware value migration through the hook (merged) |
   | PR-3 | [FLINK-40298](https://issues.apache.org/jira/browse/FLINK-40298) | 
Opt-in name-based schema evolution for `RowData` (#28973, in review) |
   | **PR-3b (this PR)** | 
[FLINK-40932](https://issues.apache.org/jira/browse/FLINK-40932) | Arm 
`ListState<RowData>` and `MapState<K, RowData>` |
   | PR-4 | [FLINK-40299](https://issues.apache.org/jira/browse/FLINK-40299) | 
End-to-end state migration coverage on RocksDB, and the user documentation |
   
   This PR is stacked on #28973 and stays a draft until it merges. Until then 
the diff below includes PR-3's two commits; only the last commit belongs to 
this step.
   
   ## What is the purpose of the change
   
   PR-3 arms schema evolution only on a state's own value serializer, so 
`ListState<RowData>` and `MapState<K, RowData>` still reject any schema change 
on restore. This PR extends the arming seam by exactly one structural level: a 
list state's element serializer and a map state's value serializer. That 
matches what the RocksDB list and map states already unwrap before they call 
`migrate`, so the invariant from PR-3 still holds: a serializer is armed only 
if some backend will invoke `migrate` on that exact serializer.
   
   Three constraints shape the change.
   
   Arming is decided by state kind, not by serializer shape. RocksDB chooses 
its migration routine by state kind: list and map states descend to the element 
or value, while value, reducing and aggregating states call `migrate` on the 
top-level serializer. A `ValueState<List<RowData>>` has the same `ListTypeInfo` 
as a `ListState<RowData>`, so a shape-based descent would arm a serializer that 
nothing migrates and reintroduce silent corruption. 
`StateSchemaEvolvingSerializer.arming` therefore takes the 
`StateDescriptor.Type`, and value, reducing and aggregating states keep exactly 
PR-3's behavior.
   
   The descent stops at one level, and never reaches a map key. A nested 
`compatibleAfterMigration` propagates up through 
`CompositeTypeSerializerSnapshot`, so arming two levels down would make the 
top-level verdict `compatibleAfterMigration` while the only `migrate` call 
lands one level up, on a snapshot that does not override it. Shapes such as the 
interval join's `MapState<Long, List<Tuple2<RowData, Boolean>>>` stay unarmed 
and fail closed. Map keys are never armed because RocksDB requires the map key 
to be `compatibleAsIs`.
   
   Operator and broadcast state now reject an armed serializer. In PR-3 they 
were excluded partly because their own serializers are a `ListSerializer` or 
`MapSerializer`, which this PR now descends into. A `StateDescriptor` caches 
its serializer per instance, so a descriptor first registered as keyed state on 
RocksDB and then reused for operator or broadcast state would carry the armed 
serializer into a backend that accepts `compatibleAfterMigration` and never 
migrates. `DefaultOperatorStateBackend` now throws a `StateMigrationException` 
for an armed serializer, using a static `isArmed` check that performs the same 
one-level descent as arming. A job that does not opt in cannot reach the check.
   
   ## Brief change log
   
     - `StateSchemaEvolvingSerializer`: kind-aware `arming(SerializerFactory, 
StateDescriptor.Type)`, one-level descent to a list element or map value, 
original composite instance returned when nothing inside it was armed, 
`isStateSchemaEvolutionEnabled()` and a static `isArmed` mirroring the descent
     - `DefaultKeyedStateStore` and `AbstractKeyedStateBackend` pass the 
descriptor's state kind to the arming factory
     - `DefaultOperatorStateBackend` rejects an armed serializer for list, 
union list and broadcast state, which also covers the state-v2 getters that 
delegate to them
     - `RowDataSerializer.isStateSchemaEvolutionEnabled()` becomes the public 
implementation of the interface method
     - The `table.exec.state.schema-evolution.enabled` description now states 
that list state elements and map state values are covered
   
   ## Verifying this change
   
   This change added tests and can be verified as follows:
   
     - A `ListState<RowData>` element and a `MapState<K, RowData>` value are 
armed, and the map key of the same descriptor is not
     - A `ValueState` whose value is a list or a map stays unarmed below the 
top level, through both `DefaultKeyedStateStore` and a direct 
`AbstractKeyedStateBackend#getOrCreateKeyedState` registration
     - The interval-join shaped `MapState<Long, List<Tuple2<RowData, 
Boolean>>>` stays unarmed and its evolved schema is rejected (existing test, 
unchanged)
     - A list or map with no schema-evolving serializer inside returns the same 
serializer instance, and a rebuilt composite is `equals` and `hashCode` equal 
to its unarmed original
     - TTL-enabled list and map states arm the element and value before TTL 
wrapping
     - A list descriptor and a map descriptor armed as keyed state are rejected 
by operator state and broadcast state respectively, and `isArmed` ignores the 
map key
   
   Each new test was checked to fail with the corresponding part of the fix 
reverted. The RocksDB state migration suites and 
`flink-architecture-tests-production` pass.
   
   ## Does this pull request potentially affect one of the following parts:
   
     - Dependencies (does it add or upgrade a dependency): no
     - The public API, i.e., is any changed class annotated with 
`@Public(Evolving)`: no. The changed interface and methods are `@Internal`.
     - The serializers: yes
     - The runtime per-record code paths (performance sensitive): no, the added 
work runs when state is registered
     - Anything that affects deployment or recovery: yes, with the option 
enabled on RocksDB, restores of list and map states holding `RowData` accept a 
backward-compatible schema change. With the option off, restore behavior is 
unchanged.
     - The S3 file system connector: no
   
   ## Documentation
   
     - Does this pull request introduce a new feature? yes, it extends the 
feature introduced in PR-3
     - If yes, how is the feature documented? The generated execution config 
option page and JavaDocs. The prose documentation lands in PR-4.
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: Claude Code (Opus 5)
   


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