João Boto created FLINK-40943:
---------------------------------

             Summary: 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


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)

Reply via email to