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

Flink窗口状态与Checkpoint故障恢复数据丢失问题咨询

Flink窗口状态与Checkpoint容错问题排查与解决

问题背景与初始疑问

现有一个用户访问数据统计任务,采用每分钟滚动窗口做求和统计,配置了间隔30秒的Checkpoint实现容错。假设任务在01:00时挂掉,我认为理论上只能恢复到00:30的状态数据,但00:30没有窗口触发,恢复后得到的是00:30的Kafka Offset数据和00:00的窗口数据,这个理解是否正确?

环境与代码说明

核心配置

  • config.getDelayMaxDuration() = 60000(水印最大乱序容忍60秒)
  • config.getAggregateWindowMillisecond()=60000(滚动窗口大小60秒)
  • Checkpoint触发间隔30秒

主程序代码

SingleOutputStreamOperator<BaseResult> wordCountSampleStream = subStream.assignTimestampsAndWatermarks(
                        WatermarkStrategy.<MetricEvent>forBoundedOutOfOrderness(config.getDelayMaxDuration())
                                .withTimestampAssigner(new MetricEventTimestampAssigner())
                                .withIdleness(config.getWindowIdlenessTime())
                ).setParallelism(CommonJobConfig.getParallelismOfSubJob("WORD_COUNT_SAMPLE_TEST"))
                .flatMap(new WordCountToResultFlatMapFunction(config)).setParallelism(CommonJobConfig.getParallelismOfSubJob("WORD_COUNT_SAMPLE_TEST"))
                .keyBy(new BaseResultKeySelector())
                .window(TumblingEventTimeWindows.of(Time.milliseconds(config.getAggregateWindowMillisecond())))
                .apply(new WordCountWindowFunction(config)).setParallelism(CommonJobConfig.getParallelismOfSubJob("WORD_COUNT_SAMPLE_TEST"));
        wordCountSampleStream.addSink(sink).setParallelism(CommonJobConfig.getParallelismOfSubJob("WORD_COUNT_SAMPLE_TEST"));

窗口Apply函数

public class WordCountWindowFunction extends RichWindowFunction<BaseResult, BaseResult, String, TimeWindow> {
private StreamingConfig config;

private Logger logger = LoggerFactory.getLogger(WordCountWindowFunction.class);

public WordCountWindowFunction(StreamingConfig config) {
    this.config = config;
}


@Override
public void close() throws Exception {
    super.close();
}

@Override
public void apply(String s, TimeWindow window, Iterable<BaseResult> input, Collector<BaseResult> out) throws Exception {
    WordCountEventPortrait result = new WordCountEventPortrait();
    long curWindowTimestamp = window.getStart() / config.getAggregateWindowMillisecond() * config.getAggregateWindowMillisecond();
    result.setDatasource("word_count_test");
    result.setTimeSlot(curWindowTimestamp);
    for (BaseResult sub : input) {
        logger.info("in window cur sub is {} ", sub);
        WordCountEventPortrait curInvoke = (WordCountEventPortrait) sub;
        result.setTotalCount(result.getTotalCount() + curInvoke.getTotalCount());
        result.setWord(curInvoke.getWord());
    }
    logger.info("out window result is {} ", result);
    out.collect(result);
}
}

Sink函数

public class ClickHouseRichSinkFunction extends RichSinkFunction<BaseResult> implements CheckpointedFunction {
private ConcurrentHashMap<String, SinkBatchInsertHelper<BaseResult>> tempResult = new ConcurrentHashMap<>();
private ClickHouseDataSource dataSource;
private Logger logger = LoggerFactory.getLogger(ClickHouseRichSinkFunction.class);


@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
    for (Map.Entry<String, SinkBatchInsertHelper<BaseResult>> helper : tempResult.entrySet()) {
        helper.getValue().insertAllTempData();
    }
}

@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
}

@Override
public void open(Configuration parameters) throws Exception {
    Properties properties = new Properties();
    properties.setProperty("user", CommonJobConfig.CLICKHOUSE_USER);
    properties.setProperty("password", CommonJobConfig.CLICKHOUSE_PASSWORD);
    dataSource = new ClickHouseDataSource(CommonJobConfig.CLICKHOUSE_JDBC_URL, properties);
}

@Override
public void close() {
    AtomicInteger totalCount = new AtomicInteger();
    tempResult.values().forEach(it -> {
        totalCount.addAndGet(it.getTempList().size());
        batchSaveBaseResult(it.getTempList());
        it.getTempList().clear();
    });
}

@Override
public void invoke(BaseResult value, Context context) {
    tempResult.compute(value.getDatasource(), (datasource, baseResults) -> {
        if (baseResults == null) {
            baseResults = new SinkBatchInsertHelper<>(CommonJobConfig.COMMON_BATCH_INSERT_COUNT,
                    needToInsert -> batchSaveBaseResult(needToInsert),
                    CommonJobConfig.BATCH_INSERT_INTERVAL_MS);
        }
        baseResults.tempInsertSingle(value);
        return baseResults;
    });
}

private void batchSaveBaseResult(List<BaseResult> list) {
    if (list.isEmpty()) {
        return;
    }
    String sql = list.get(0).getPreparedSQL();
    try {
        try (PreparedStatement ps = dataSource.getConnection().prepareStatement(sql)) {
            for (BaseResult curResult : list) {
                curResult.addParamsToPreparedStatement(ps);
                ps.addBatch();
            }
            ps.executeBatch();
        }
    } catch (SQLException error) {
        logger.error("has exception during batch insert,datasource is {} ", list.get(0).getDatasource(), error);
    }
}
}

批量插入工具类

public class SinkBatchInsertHelper<T> {
private List<T> waitToInsert;

private ReentrantLock lock;

private int bulkActions;

private AtomicInteger tempActions;

private Consumer<List<T>> consumer;

private AtomicLong lastSendTimestamp;

private long sendInterval;

private Logger logger = LoggerFactory.getLogger(SinkBatchInsertHelper.class);

public SinkBatchInsertHelper(int bulkActions, Consumer<List<T>> consumer, long sendInterval) {
    this.waitToInsert = new ArrayList<>();
    this.lock = new ReentrantLock();
    this.bulkActions = bulkActions;
    this.tempActions = new AtomicInteger(0);
    this.consumer = consumer;
    this.sendInterval = sendInterval;
    this.lastSendTimestamp = new AtomicLong(0);
}


public void tempInsertSingle(T data) {
    lock.lock();
    try {
        waitToInsert.add(data);
        if (tempActions.incrementAndGet() >= bulkActions || ((System.currentTimeMillis() - lastSendTimestamp.get()) >= sendInterval)) {
            batchInsert();
        }
    } finally {
        lastSendTimestamp.set(System.currentTimeMillis());
        lock.unlock();
    }
}

public long insertAllTempData() {
    lock.lock();
    try {
        long result = tempActions.get();
        if (tempActions.get() > 0) {
            batchInsert();
        }
        return result;
    } finally {
        lock.unlock();
    }
}

private void batchInsert() {
    for(T t: waitToInsert){
            logger.info("batch insert data:{}", t);
    }
    consumer.accept(waitToInsert);
    waitToInsert.clear();
    tempActions.set(0);
}


public int getTempActions() {
    return tempActions.get();
}

public List<T> getTempList() {
    lock.lock();
    try {
        return waitToInsert;
    } finally {
        lock.unlock();
    }
}
}

实际异常现象

若在00:31:30取消任务,重启后00:31:00的统计数据低于预期。日志显示:Sink写入00:30:00的数据时,Kafka消费者已消费00:31:00之后的数据,但这部分数据未写入Checkpoint,重启后未在窗口中重放,导致数据丢失,直到00:32:00统计才恢复正常。

问题解答

1. 初始理解的正确性

你的理解不完全正确,核心误区在于对Checkpoint状态保存逻辑的误解:

  • Checkpoint是增量覆盖式的,每次成功完成的Checkpoint会保存当前所有算子的最新状态,包括:Kafka Consumer的最新Offset、窗口的中间聚合状态、已触发窗口的处理记录。
  • 若任务在01:00挂掉,恢复时会回到最近一次成功完成的Checkpoint状态(比如00:59:30触发的Checkpoint):
    • Kafka Offset会恢复到00:59:30时消费到的位置;
    • 窗口状态会包含00:59:00-01:00:00窗口的中间聚合数据(如果已有数据进入);
    • 00:59:00之前的窗口(如00:58:00-00:59:00)若已触发完成,其状态会被清理,对应的结果如果已通过Sink持久化,Checkpoint会保证数据要么已写入,要么会被重放。

2. 数据丢失问题的根源与解决方法

根源分析

当前自定义Sink的实现破坏了Flink Checkpoint的Exactly-Once语义:

  1. Checkpoint阶段直接写入外部存储:在snapshotState中直接调用insertAllTempData()写入ClickHouse,若Checkpoint完成前任务挂掉,已写入的数据无法回滚,而Kafka Offset已被保存,重启后会跳过这部分数据;若写入失败,Checkpoint会失败,但已消费的数据状态已更新,导致数据不一致。
  2. 批量写入与Checkpoint的时序不匹配:Sink的批量写入逻辑(按数量或时间触发)与Checkpoint触发时机无关联,可能出现数据已写入但未被Checkpoint记录,或数据未写入但Offset已被保存的情况。

解决方法

方案一:使用官方ClickHouse Connector

Flink 1.14已提供官方ClickHouse Connector,内置两阶段提交(2PC)机制,原生支持Exactly-Once语义,无需手动实现复杂的Checkpoint逻辑,直接替换自定义Sink即可。

方案二:修复自定义Sink的Exactly-Once语义

若必须使用自定义Sink,需实现两阶段提交逻辑,核心步骤如下:

  1. 在snapshotState中,将待提交的数据缓存到Flink的算子状态中,而非直接写入外部存储;
  2. 实现CheckpointListener接口,在notifyCheckpointComplete回调中,将缓存的数据批量写入ClickHouse;
  3. 在initializeState中恢复故障前的缓存数据,根据Checkpoint完成状态决定是否提交或回滚;
  4. 确保批量写入逻辑仅在Checkpoint完成后触发,避免数据提前提交。

修复后的Sink代码示例:

public class ClickHouseRichSinkFunction extends RichSinkFunction<BaseResult> implements CheckpointedFunction, CheckpointListener {
    private transient ListState<BaseResult> checkpointedState;
    private List<BaseResult> pendingRecords = new ArrayList<>();
    private ClickHouseDataSource dataSource;
    private Logger logger = LoggerFactory.getLogger(ClickHouseRichSinkFunction.class);

    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        checkpointedState.clear();
        for (BaseResult record : pendingRecords) {
            checkpointedState.add(record);
        }
    }

    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        ListStateDescriptor<BaseResult> descriptor = new ListStateDescriptor<>(
                "pending-records",
                TypeInformation.of(BaseResult.class)
        );
        checkpointedState = context.getOperatorStateStore().getListState(descriptor);

        if (context.isRestored()) {
            for (BaseResult record : checkpointedState.get()) {
                pendingRecords.add(record);
            }
        }
    }

    @Override
    public void notifyCheckpointComplete(long checkpointId) throws Exception {
        if (!pendingRecords.isEmpty()) {
            batchSaveBaseResult(pendingRecords);
            pendingRecords.clear();
            checkpointedState.clear();
        }
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        Properties properties = new Properties();
        properties.setProperty("user", CommonJobConfig.CLICKHOUSE_USER);
        properties.setProperty("password", CommonJobConfig.CLICKHOUSE_PASSWORD);
        dataSource = new ClickHouseDataSource(CommonJobConfig.CLICKHOUSE_JDBC_URL, properties);
    }

    @Override
    public void invoke(BaseResult value, Context context) {
        pendingRecords.add(value);
    }

    private void batchSaveBaseResult(List<BaseResult> list) {
        if (list.isEmpty()) {
            return;
        }
        String sql = list.get(0).getPreparedSQL();
        try (PreparedStatement ps = dataSource.getConnection().prepareStatement(sql)) {
            for (BaseResult curResult : list) {
                curResult.addParamsToPreparedStatement(ps);
                ps.addBatch();
            }
            ps.executeBatch();
            logger.info("批量写入{}条数据成功", list.size());
        } catch (SQLException error) {
            logger.error("批量插入失败,数据源:{}", list.get(0).getDatasource(), error);
            throw new RuntimeException("批量插入失败", error);
        }
    }

    @Override
    public void close() throws Exception {
        if (!pendingRecords.isEmpty()) {
            batchSaveBaseResult(pendingRecords);
        }
    }
}
方案三:调整Checkpoint与窗口的配合
  • 将Checkpoint间隔调整为与窗口大小一致(60秒),确保每次Checkpoint在窗口触发完成后执行,减少中间状态不一致的概率,但会降低容错恢复的粒度;
  • 调整水印乱序容忍时间:当前设置的60秒与窗口大小一致,会导致窗口触发延迟60秒,可根据实际数据乱序程度缩小该值,避免不必要的延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 22:00:55