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语义:
- Checkpoint阶段直接写入外部存储:在
snapshotState中直接调用insertAllTempData()写入ClickHouse,若Checkpoint完成前任务挂掉,已写入的数据无法回滚,而Kafka Offset已被保存,重启后会跳过这部分数据;若写入失败,Checkpoint会失败,但已消费的数据状态已更新,导致数据不一致。 - 批量写入与Checkpoint的时序不匹配:Sink的批量写入逻辑(按数量或时间触发)与Checkpoint触发时机无关联,可能出现数据已写入但未被Checkpoint记录,或数据未写入但Offset已被保存的情况。
解决方法
方案一:使用官方ClickHouse Connector
Flink 1.14已提供官方ClickHouse Connector,内置两阶段提交(2PC)机制,原生支持Exactly-Once语义,无需手动实现复杂的Checkpoint逻辑,直接替换自定义Sink即可。
方案二:修复自定义Sink的Exactly-Once语义
若必须使用自定义Sink,需实现两阶段提交逻辑,核心步骤如下:
- 在
snapshotState中,将待提交的数据缓存到Flink的算子状态中,而非直接写入外部存储; - 实现
CheckpointListener接口,在notifyCheckpointComplete回调中,将缓存的数据批量写入ClickHouse; - 在
initializeState中恢复故障前的缓存数据,根据Checkpoint完成状态决定是否提交或回滚; - 确保批量写入逻辑仅在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
相关产品推荐
相关产品推荐

