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]
