Skip to content

[FLINK-40943][filesystems] Restore FileSink state written with flink-s3-fs-hadoop in flink-s3-fs-native - #29422

Open
eskabetxe wants to merge 1 commit into
apache:masterfrom
eskabetxe:FLINK-40943_filesink_restore_state
Open

eskabetxe wants to merge 1 commit into
apache:masterfrom
eskabetxe:FLINK-40943_filesink_restore_state

Conversation

@eskabetxe

@eskabetxe eskabetxe commented Oct 7, 2026 •

Copy link
Copy Markdown
Member

What is the purpose of the change

Fixes 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?
  • [] Yes (please specify the tool below)

@eskabetxe
eskabetxe force-pushed the FLINK-40943_filesink_restore_state branch from 99bd15e to b255687 Compare October 7, 2026 13:59
@flinkbot

flinkbot commented Oct 7, 2026 •

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

@gaborgsomogyi

Copy link
Copy Markdown
Contributor

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?

@eskabetxe

Copy link
Copy Markdown
Member Author

Hi @gaborgsomogyi, thanks for looking!

Checkpoint vs savepoint: yes, identical behavior. Both store the same FileSinkCommittable bytes through the same serializer chain (FileSinkCommittableSerializer -> OutputStreamBasedPartFileWriter recoverable serializers -> the filesystem's RecoverableWriter serializer); a savepoint is just a portable encoding of the same bytes. So restoring either with flink-s3-fs-native hits the same legacy-decode branch in CommitterOperator#initializeState — that's also where the original EOFException came from.

Endianness: nothing is hardcoded ad-hoc — each format's endianness comes from the API that writes it. The legacy S3RecoverableSerializer (flink-s3-fs-base) writes via ByteBuffer.order(LITTLE_ENDIAN) (see its serialize), so the new decoder reads it back with ByteBuffer.wrap(serialized).order(LITTLE_ENDIAN) in isLegacyFormat/decodeLegacy. The native format uses DataOutputStream/DataInputStream, which the Java spec defines as big-endian. Detection is unambiguous in practice: mistaking native data for the legacy magic would require the first writeUTF length (big-endian, i.e. the object name's length) to be 0x3214 = 12820 bytes, far beyond the 1024-byte S3 key limit.

@gaborgsomogyi gaborgsomogyi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +44 to +46
// 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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Comment on lines +156 to +158
{{< 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 >}}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

@eskabetxe eskabetxe Oct 9, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed.
But same comment as before: if we are left behind the old formats, they should not receive more changes.

Comment on lines +178 to +184
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;
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Under what circumstances can the legacy format magic collide with a valid new format checkpoint?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This already exists

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

changed to use the value from the constant on the other class.

@eskabetxe
eskabetxe force-pushed the FLINK-40943_filesink_restore_state branch from b255687 to 566314a Compare October 9, 2026 13:18
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants