Flink 1.15.1 Table API流式S3 Sink触发SerializedThrowable异常求助
Flink 1.15.1 Table API S3 CSV Sink Checkpoint触发异常排查方案
问题场景
使用Flink 1.15.1的Table API实现流式数据写入S3(s3a协议),格式为CSV。触发Checkpoint时必现如下异常,更换Kafka/datagen数据源均无法解决,仅在Checkpoint阶段失败。
Caused by: org.apache.flink.util.SerializedThrowable: S3RecoverableFsDataOutputStream cannot sync state to S3. Use persist() to create a persistent recoverable intermediate point. at org.apache.flink.fs.s3.common.utils.RefCountedBufferingFileStream.sync(RefCountedBufferingFileStream.java:111) ~[flink-s3-fs-hadoop-1.15.1.jar:1.15.1] at org.apache.flink.fs.s3.common.writer.S3RecoverableFsDataOutputStream.sync(S3RecoverableFsDataOutputStream.java:129) ~[flink-s3-fs-hadoop-1.15.1.jar:1.15.1] at org.apache.flink.formats.csv.CsvBulkWriter.finish(CsvBulkWriter.java:110) ~[flink-csv-1.15.1.jar:1.15.1] at org.apache.flink.connector.file.table.FileSystemTableSink$ProjectionBulkFactory$1.finish(FileSystemTableSink.java:642) ~[flink-connector-files-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:64) ~[flink-file-sink-common-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.closePartFile(Bucket.java:263) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.prepareBucketForCheckpointing(Bucket.java:305) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.onReceptionOfCheckpoint(Bucket.java:277) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.Buckets.snapshotActiveBuckets(Buckets.java:270) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.Buckets.snapshotState(Buckets.java:261) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSinkHelper.snapshotState(StreamingFileSinkHelper.java:87) ~[flink-streaming-java-1.15.1.jar:1.15.1] at org.apache.flink.connector.file.table.stream.AbstractStreamingWriter.snapshotState(AbstractStreamingWriter.java:129) ~[flink-connector-files-1.15.1.jar:1.15.1]
异常原因
核心矛盾在于CSV BulkWriter的finish方法会强制调用文件流的sync操作,但S3作为对象存储,不支持本地文件系统的sync语义。Flink的S3文件系统实现要求通过persist()创建可恢复的持久化中间点,而非sync(),两者行为不兼容导致Checkpoint时失败。
解决方案
1. 切换至对象存储友好的文件格式(推荐)
改用Parquet、ORC等列式存储格式,这类格式的BulkWriter实现适配了S3的特性,不会触发sync调用。示例SQL:
CREATE TABLE s3_sink ( id INT, name STRING, dt STRING ) WITH ( 'connector' = 'filesystem', 'path' = 's3a://your-bucket/target-path', 'format' = 'parquet', 'partition.column' = 'dt', 'sink.partition-commit.trigger' = 'checkpoint', 'sink.partition-commit.policy.kind' = 'success-file' );
2. 配置S3A专属提交器(CSV格式适配)
若必须使用CSV格式,需为S3A配置适配的提交器,避免sync操作冲突。
- 在
flink-conf.yaml中添加:fs.s3a.committer.factory.class=org.apache.flink.fs.s3a.commit.S3ACommitterFactory fs.s3a.committer.name=partitioned fs.s3a.committer.staging.dir=.tmp/staging fs.s3a.committer.staging.conflict-mode=replace - 同时在Table定义中指定分区提交策略:
CREATE TABLE s3_csv_sink ( id INT, name STRING, dt STRING ) WITH ( 'connector' = 'filesystem', 'path' = 's3a://your-bucket/csv-path', 'format' = 'csv', 'partition.column' = 'dt', 'sink.partition-commit.trigger' = 'checkpoint', 'sink.partition-commit.policy.kind' = 'success-file' );
3. 排查依赖冲突
检查项目依赖树,确保未同时引入flink-s3-fs-hadoop和flink-s3-fs-presto两个依赖。若Checkpoint使用s3p协议、Sink使用s3a协议,可能导致文件系统实现类加载冲突,建议统一使用一种S3协议(如全量用s3a),或调整依赖排除冲突。
内容的提问来源于stack exchange,提问作者Jaya Ananthram
相关产品推荐
相关产品推荐

