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

如何配置避免Flink处理Kinesis数据流时丢失故障记录?

解决Flink处理Kinesis记录时的丢失问题

1. 启用Flink检查点并配置精准一次语义

必须开启检查点机制,让Flink能在故障恢复时回滚到一致状态,并重播未完成处理的记录:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 开启检查点,间隔可根据业务调整示例为5分钟
env.enableCheckpointing(300000);
// 配置精准一次语义,确保记录仅被处理一次
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 设置检查点超时时间,避免长时间阻塞
env.getCheckpointConfig().setCheckpointTimeout(600000);
// 检查点失败时终止作业,防止状态不一致
env.getCheckpointConfig().setFailOnCheckpointingErrors(true);
// 保留外部检查点,用于作业终止后的恢复
env.getCheckpointConfig().enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

该配置会保证Kinesis偏移量仅在记录被成功处理并写入Sink后,才会随检查点提交,故障恢复时会从上次成功的偏移量重新消费,不会丢失记录。

2. 配置FlinkKinesisConsumer的偏移量提交策略

禁用KCL自动提交偏移量,改为由Flink检查点控制提交时机:

Properties kinesisProps = new Properties();
kinesisProps.setProperty(ConsumerConfigConstants.STREAM_INITIAL_POSITION, "LATEST");
kinesisProps.setProperty(ConsumerConfigConstants.SHARD_GETRECORDS_INTERVAL_MILLIS, "1000");
// 关闭定时自动提交偏移量
kinesisProps.setProperty(ConsumerConfigConstants.COMMIT_INTERVAL_MS, "-1");
// 启用检查点驱动的偏移量提交
kinesisProps.setProperty(ConsumerConfigConstants.USE_CHECKPOINTED_OFFSET, "true");

FlinkKinesisConsumer<String> kinesisConsumer = new FlinkKinesisConsumer<>(
    "your-stream-name",
    new SimpleStringSchema(),
    kinesisProps
);
env.addSource(kinesisConsumer);

此配置确保只有当记录被成功处理并完成检查点后,偏移量才会被提交,避免处理失败后偏移量已提交导致的记录丢失。

3. 为MapFunction配置重试与异常处理

配置重启策略

根据业务需求设置重启策略,确保瞬时异常能通过重试恢复,同时避免无意义的循环重启:

// 固定延迟重试:最多重试3次,每次间隔10秒
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
    3, 
    Time.seconds(10)
));

// 若需要无限重试直到手动干预,可使用:
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(Integer.MAX_VALUE, Time.seconds(30)));

添加死信队列处理不可恢复异常

对于持续性异常,可捕获异常并将无法处理的记录发送到死信队列,避免作业循环崩溃同时留存记录:

public class MyMapFunction extends MapFunction<String, String> {
    private transient OutputTag<String> deadLetterOutputTag;

    @Override
    public void open(Configuration parameters) throws Exception {
        deadLetterOutputTag = new OutputTag<String>("dead-letter-records"){};
    }

    @Override
    public String map(String value) throws Exception {
        try {
            // 业务处理逻辑
            return processValue(value);
        } catch (Exception e) {
            // 将异常记录发送到死信流
            getRuntimeContext().output(deadLetterOutputTag, value);
            // 返回空值或标记值,后续可过滤,避免抛出异常导致作业重启
            return null;
        }
    }
}

// 主作业中处理死信流
SingleOutputStreamOperator<String> mainStream = env.addSource(kinesisConsumer)
    .map(new MyMapFunction());

// 将死信记录写入专门的Sink
mainStream.getSideOutput(new OutputTag<String>("dead-letter-records"){})
    .addSink(new DeadLetterSink());

4. 禁用KCL自动跳过记录

确保不配置KCL中可能导致自动跳过记录的参数,同时依赖Flink检查点控制偏移量提交,即可避免KCL静默跳过记录的行为。

效果总结

通过以上配置可实现:

  • 记录处理失败时,偏移量不会被提交,恢复时会重处理该记录
  • 瞬时异常可通过重试策略恢复,不会丢失记录
  • 持续性异常可通过死信队列留存记录,避免作业循环崩溃同时保证无数据丢失

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 19:06:28