Flink savepoint compatibility mode behavior after a breaking schema change – rollback safety?
0 reputation · 30 Sept 2024, 11:53 UTC
0 reputation · 30 Sept 2024, 11:53 UTC
The goal is to assess if a Flink job can be safely rolled back to a pre‑change schema after a breaking modification (e.g., field removal or type change) when using the savepoint compatibility mode introduced in Flink 1.16.
Current documentation states that automatic restoration is only possible for backward‑compatible schema evolutions; for breaking changes the savepoint must be manually transformed or the job restarted from scratch, and the State Processor API does not provide transactional consistency guarantees.
Given these constraints, what verified procedures exist to guarantee rollback safety without risking data corruption or deserialization failures?
29775 reputation · 30 Sept 2024, 23:36 UTC
In Flink 1.18+ the compatibility mode can only auto‑restore for backward‑compatible changes (e.g., adding fields or widening types). When a field is removed or its type is narrowed, the job must either keep the original savepoint and restore from it, or run a custom state migration before restarting the job. The compatibility mode does not provide automatic rollback for breaking changes.
state.checkpointing.mode and state.backend.checkpointing.mode but does not affect breaking schema changes.flink savepoint --target-directory hdfs://... --savepoint <job-id> /path/to/savepointStateTransformer that reads the old state format and writes the new one (or vice‑versa). Package it into a JAR.flink run -c org.apache.flink.stateprocessor.StateProcessorJob /path/to/transformer.jar \
--checkpoint-dir hdfs://... --target-dir hdfs://... \
--source-savepoint /path/to/savepoint --target-savepoint /path/to/migratedflink run -c com.example.NewJob /path/to/new-job.jar --savepoint /path/to/migratedflink run -c com.example.OldJob /path/to/old-job.jar --savepoint /path/to/savepointflink state-tool --metadata on both the original and migrated savepoints to verify state size and structure.rocksdb-check after migration to detect corruption.To tailor the migration script precisely, could you confirm whether your state backend is RocksDB or FsStateBackend? This choice can affect how the State Processor API reads and writes the state files.
Use comments to ask for clarification. Post a solution as an answer.
29,775 reputation · 30 Sept 2024, 14:29 UTC
The execution.savepoint.compatibility-mode setting only relaxes Flink runtime version constraints (e.g., restoring a 1.17 savepoint on 1.18); it does not alter schema evolution rules. Breaking changes still require a manual State Processor migration.
For validation, the community pattern is a reversible State Processor job run on a shadow cluster: read the pre-change savepoint, apply the forward migration, then immediately run the reverse migration and compare key/value counts and RocksDB SST checksums against the original. Only promote after dual-run output equivalence is confirmed on a bounded test dataset.
Caveat: if the pre-change savepoint was incremental (RocksDB), ensure it contains all referenced segment files — some older segments may have been cleaned by the retention policy, making full rollback impossible without a prior full checkpoint.