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]

Reply via email to