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

如何将Flink S3 Sink的投递模式从EXACTLY_ONCE改为AT_LEAST_ONCE

核心结论

作业级的Checkpoint模式修改不会自动同步更新S3 Sink的投递模式,Sink的投递保障是独立配置的,需要显式调整。

如何修改S3 Sink为AT_LEAST_ONCE模式

Flink的S3 Sink(如FileSink)通过DeliveryGuarantee参数控制投递模式,默认是EXACTLY_ONCE,需在构建Sink时显式设置为AT_LEAST_ONCE:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.core.fs.Path;
import org.apache.flink.formats.string.SimpleStringEncoder;
import org.apache.flink.streaming.api.functions.sink.filesystem.DeliveryGuarantee;
import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.DateTimeBucketAssigner;
import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy;

// 构建AT_LEAST_ONCE模式的S3 Sink
FileSink<String> s3Sink = FileSink
    .forRowFormat(new Path("s3://your-bucket/target-path"), new SimpleStringEncoder<String>("UTF-8"))
    .withRollingPolicy(DefaultRollingPolicy.builder().build())
    .withBucketAssigner(new DateTimeBucketAssigner<>("yyyy-MM-dd"))
    .withDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) // 关键配置:显式设置投递模式
    .build();

// 将Sink添加到作业
dataStream.sinkTo(s3Sink);

Checkpoint模式与Sink投递模式的关系

  1. 作业的CheckpointingMode决定整个作业的一致性语义,而Sink的DeliveryGuarantee是该Sink自身的投递保障策略,二者关联但独立:

    • 作业设为EXACTLY_ONCE时,Sink可选EXACTLY_ONCE(配合Checkpoint做原子提交)或AT_LEAST_ONCE(跳过原子提交直接写入)
    • 作业设为AT_LEAST_ONCE时,Sink必须设为AT_LEAST_ONCE,否则作业启动失败(EXACTLY_ONCE的Sink依赖作业的精准一次Checkpoint机制)
  2. 若你将作业Checkpoint模式改为AT_LEAST_ONCE:

    env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);
    

    必须同步修改S3 Sink的DeliveryGuarantee为AT_LEAST_ONCE,否则会抛出兼容性异常。

模式切换的影响

AT_LEAST_ONCE模式下,S3 Sink会直接写入最终文件,不再使用「临时文件+Checkpoint完成后重命名」的原子提交机制,可能产生重复数据,但吞吐量会比EXACTLY_ONCE更高。

内容的提问来源于stack exchange,提问作者priyadhingra19

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 00:00:17