Repository navigation
Conversation
99bd15e to
b255687
Compare
|
Hey, thanks for filing it. Just for my own understanding is this the same behavior for checkpoint and savepoint too? Additionally what is the magic in the code what hardcodes endian? |
|
Hi @gaborgsomogyi, thanks for looking! Checkpoint vs savepoint: yes, identical behavior. Both store the same Endianness: nothing is hardcoded ad-hoc — each format's endianness comes from the API that writes it. The legacy |
gaborgsomogyi
left a comment
There was a problem hiding this comment.
As a general saying. Almost everybody is using AI to generate code but seeing complete answers doesn't strengthen the trust. We have agents locally so we can use our own words/ideas here.
| // The wire format below is frozen: flink-s3-fs-native (NativeS3RecoverableSerializer) | ||
| // duplicates this decoder to restore state after switching to that plugin (FLINK-40943). | ||
| // Do not change the format without keeping that decoder in sync. |
There was a problem hiding this comment.
I would say this support with the copy-paste is one time step and any further enhancements in the serializer logic will break this contract. To be specific we give users a one time path to the new connector
There was a problem hiding this comment.
Could be one time step, if the old serializer ever changed, we would have to copy the new format again.
I dont know the roadmap, but if we are left behind the old formats, they should not receive more changes, correct?
| {{< hint info >}} | ||
| **Switching from the Hadoop or Presto implementations:** jobs with pending or in-progress `FileSink` part files in their state can be restored directly with the Native S3 FileSystem — checkpoints and savepoints written by the Hadoop/Presto recoverable-state format are detected and migrated on restore. The reverse direction is **not** supported: state written by the Native S3 FileSystem cannot be restored with the Hadoop or Presto implementations. Drain the job first (e.g. `stop-with-savepoint --drain`) before switching back. | ||
| {{< /hint >}} |
There was a problem hiding this comment.
Two things to mention:
- We're going to leave the old connectors behind
- This is a one time occasion to have a migration path. Any further development in the old connectors will not be reflected in the native connector
There was a problem hiding this comment.
Changed.
But same comment as before: if we are left behind the old formats, they should not receive more changes.
| private static boolean isLegacyFormat(byte[] serialized) { | ||
| if (serialized.length < Integer.BYTES) { | ||
| return false; | ||
| } | ||
| return ByteBuffer.wrap(serialized).order(ByteOrder.LITTLE_ENDIAN).getInt() | ||
| == LEGACY_MAGIC_NUMBER; | ||
| } |
There was a problem hiding this comment.
Under what circumstances can the legacy format magic collide with a valid new format checkpoint?
There was a problem hiding this comment.
Practically none.
The first 2 bytes would only be achieved wit object names with 12820 bytes and S3 keys cap at 1024.
Also 4 byte is 0x98 witch is not a valid UTF-8 character byte
| /** Tests for {@link NativeS3RecoverableSerializer}. */ | ||
| class NativeS3RecoverableSerializerTest { | ||
|
|
||
| private static final int LEGACY_MAGIC_NUMBER = 0x98761432; |
There was a problem hiding this comment.
changed to use the value from the constant on the other class.
…s3-fs-hadoop in flink-s3-fs-native
b255687 to
566314a
Compare
What is the purpose of the change
Fixes FLINK-40943.
Restoring a FileSink job from a checkpoint/savepoint taken with
flink-s3-fs-hadoop(orflink-s3-fs-presto) fails after switching toflink-s3-fs-native, becauseNativeS3RecoverableSerializerandS3RecoverableSerializer(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 crypticEOFExceptioninCommitterOperator#initializeState, leaving no migration path other than dropping state.This PR makes
NativeS3RecoverableSerializer.deserializedetect the legacy format via its leading magic number (0x98761432) and decode it intoNativeS3Recoverable, so jobs with pending or in-progress part files can be restored directly after switching toflink-s3-fs-native.Note: compatibility is one-way. State written by
flink-s3-fs-nativestill cannot be restored withflink-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 legacyS3RecoverableSerializerformat (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 descriptiveIOException. 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.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:S3RecoverableSerializer#serializein 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.IOException.Does this pull request potentially affect one of the following parts:
@Public(Evolving): noDocumentation
deployment/filesystems/s3.md)Was generative AI tooling used to co-author this PR?