eskabetxe opened a new pull request, #29422: URL: https://github.com/apache/flink/pull/29422
## What is the purpose of the change Fixes [FLINK-40943](https://issues.apache.org/jira/browse/FLINK-40943). Restoring a FileSink job from a checkpoint/savepoint taken with `flink-s3-fs-hadoop` (or `flink-s3-fs-presto`) fails after switching to `flink-s3-fs-native`, because `NativeS3RecoverableSerializer` and `S3RecoverableSerializer` (flink-s3-fs-base) encode the multipart-upload recoverable state differently but both report serializer version 1. The native reader misinterprets the leading magic number of the legacy format and fails with a cryptic `EOFException` in `CommitterOperator#initializeState`, leaving no migration path other than dropping state. This PR makes `NativeS3RecoverableSerializer.deserialize` detect the legacy format via its leading magic number (`0x98761432`) and decode it into `NativeS3Recoverable`, so jobs with pending or in-progress part files can be restored directly after switching to `flink-s3-fs-native`. Note: compatibility is one-way. State written by `flink-s3-fs-native` still cannot be restored with `flink-s3-fs-hadoop`/`presto`; jobs must be drained (e.g. `stop-with-savepoint --drain`) before switching back. This is now stated in the docs. ## Brief change log - `NativeS3RecoverableSerializer` (flink-s3-fs-native): detect the legacy `S3RecoverableSerializer` format (flink-s3-fs-base) via its magic number and decode it (field-for-field identical semantics; both recoverables share the same invariants). Length fields and part counts are validated against the remaining buffer, trailing bytes are rejected, and corrupt input is wrapped in a descriptive `IOException`. An INFO log is emitted when legacy state is detected. - `S3RecoverableSerializer` (flink-s3-fs-base): comment marking the wire format as frozen because flink-s3-fs-native duplicates its decoder. - Docs: note on the s3-native page that hadoop/presto -> native restores work, and native -> hadoop/presto requires draining first. The legacy decoder is duplicated from flink-s3-fs-base rather than shared: the two plugins share no module where the decoder would fit (different AWS SDK versions/types), are shaded and loaded in isolated plugin classloaders, and the base decoder must stay untouched for backportability. ## Verifying this change Added tests in `NativeS3RecoverableSerializerTest`: - Golden test against bytes captured from the real `S3RecoverableSerializer#serialize` in flink-s3-fs-base (parts + incomplete part + multi-byte/non-ASCII keys, incl. a supplementary character), pinning the legacy wire format independently of the test-side re-encoder. - Legacy decode tests (with/without incomplete part, empty parts, multi-byte UTF-8 names) using a test-side re-encoder whose output was verified byte-identical against the real flink-s3-fs-base serializer. - Corrupt-input tests: truncated buffer, trailing bytes, negative/huge field length, huge part count — all assert a descriptive `IOException`. - Current-format round trips and a test that current-format data is never mistaken for the legacy format. ## 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 serializers: yes (adds decoding of the legacy S3 recoverable-state format in flink-s3-fs-native; no change to any serialized output) - The runtime per-record code paths (performance sensitive): no - Anything that affects deployment or recovery: JobManager (and its components), Checkpointing, Kubernetes/Yarn, ZooKeeper: yes (checkpoint/savepoint restore of FileSink state on S3; restores that previously failed now succeed) - The S3 file system connector: yes ## Documentation - Does this pull request introduce a new feature? no - If yes, how is the feature documented? docs (migration note added to `deployment/filesystems/s3.md`) --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes (please specify the tool below) Generated-by: opencode 1.18.34 -- 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]
