[ 
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)

Reply via email to