[
https://issues.apache.org/jira/browse/FLINK-40943?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40943:
-----------------------------------
Labels: pull-request-available (was: )
> FileSink state written with flink-s3-fs-hadoop cannot be restored with
> flink-s3-fs-native (incompatible recoverable serializers, both version 1)
> ------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: FLINK-40943
> URL: https://issues.apache.org/jira/browse/FLINK-40943
> Project: Flink
> Issue Type: Bug
> Components: Connectors / FileSystem
> Affects Versions: 2.3.0
> Reporter: João Boto
> Priority: Major
> Labels: pull-request-available
>
> h3. Problem
> Restoring a job that uses FileSink to S3 from a checkpoint/savepoint taken
> with flink-s3-fs-hadoop, after switching to flink-s3-fs-native, fails during
> CommitterOperator#initializeState with an EOFException that doesn't say what
> went wrong:
>
> {code:java}
> org.apache.flink.util.FlinkRuntimeException: Failed to deserialize value
> at
> ...SimpleVersionedListState$DeserializingIterator.next(SimpleVersionedListState.java:140)
> at ...sink.CommitterOperator.initializeState(CommitterOperator.java:136)
> Caused by: java.io.EOFException
> at java.io.DataInputStream.readUTF(...)
> at
> org.apache.flink.fs.s3native.writer.NativeS3RecoverableSerializer.deserialize(NativeS3RecoverableSerializer.java:113)
> at
> ...OutputStreamBasedPartFileWriter$OutputStreamBasedPendingFileRecoverableSerializer.deserializeV2(OutputStreamBasedPartFileWriter.java:537)
> at
> ...file.sink.FileSinkCommittableSerializer.deserializeV2(FileSinkCommittableSerializer.java:138){code}
>
> h3. Cause
> The two serializers use different formats but both report getVersion() == 1:
> * S3RecoverableSerializer (flink-s3-fs-base / hadoop): little-endian
> ByteBuffer starting with magic number 0x98761432, then length-prefixed byte
> arrays.
> * NativeS3RecoverableSerializer: DataOutputStream
> writeUTF/writeLong/writeInt sequence.
> Since the versions match, the version check passes and the native reader
> interprets the first two bytes of the legacy magic number (0x32 0x14, read as
> a big-endian length of 12820) as the length of the objectName UTF string,
> reads past the end of the buffer and fails.
> h3. Impact
> Any job that has pending or in-progress part files in its state when the
> filesystem is switched (in practice, any actively writing FileSink job)
> cannot restore. The only ways out are to drop the state, which loses the
> committed-but-unpublished files and leaves orphaned multipart uploads, or to
> roll back to the old filesystem.
> h3. Suggested fix (either)
> * Read the legacy format in NativeS3RecoverableSerializer: detect the
> 0x98761432 magic number and decode the S3Recoverable layout (same fields:
> object name, upload id, part ETags, bytes in parts, incomplete object
> name/length).
> * At minimum, fail fast with a clear error (e.g. "state was written by
> flink-s3-fs-hadoop; drain the job with stop-with-savepoint --drain before
> switching"), and document that migration step on the native S3 filesystem
> page.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)