You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.15.1 Table API流式S3 Sink触发SerializedThrowable异常求助

问题场景

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 06:27:28