Apache Paimon 数据出仓源码导读(七):任务失败后数据怎样重来:Checkpoint、重试与幂等恢复

发布时间:2026/9/4 6:02:45
Apache Paimon 数据出仓源码导读(七):任务失败后数据怎样重来:Checkpoint、重试与幂等恢复 上一篇我们拆开了 MySQL UPSERT、DELETE 和 Flush 边界。这一篇专门讨论失败。我们会把任务停在最容易误判的位置JDBC Batch 已经成功写入 MySQL 但当前 Checkpoint 还没有全局成功这时 TaskManager 故障作业恢复后同一批数据会不会再来一次答案是普通 Flink SQL JDBC Sink 允许重放属于 At-Least-Once它不保证每条 DML 只执行一次而是依靠目标表主键 UPSERT 和 DELETE 让重复执行收敛到相同的最终状态。下面把这句话逐层拆开。一、先区分三个不同的“成功”初学者最容易把下面三件事都叫作“Checkpoint 成功”时刻实际含义MySQL 是否可能已看到数据故障后是否可能重放JDBCexecuteBatch()成功当前数据库 Batch 执行完成是可能SinksnapshotState()返回当前 Sink 子任务完成快照动作是仍可能Checkpoint 全局完成Source、算子、Sink 等全部确认是它才成为新的稳定恢复点核心时间差是MySQL 可以很早就看到数据 Flink 要等所有参与者都完成才拥有新的全局恢复点二、Checkpoint 保存的到底是什么对这条链路来说Checkpoint 需要协调多类状态Paimon Source 已规划哪些 Snapshot 哪些 Split 已经分配 Split 读到什么位置 中间算子 Keyed State、Operator State、Timer 等 JDBC Sink 在快照前把内存 Buffer Flush 出去Checkpoint 不是简单地“给所有线程拍一张照片”。它要让不同并行任务对同一个 Checkpoint ID 达成一致并由 JobManager 确认全局完成。三、为什么 JDBC Sink 在snapshotState()中 FlushGenericJdbcSinkFunction的关键代码publicvoidsnapshotState(FunctionSnapshotContextcontext)throwsException{outputFormat.flush();}为什么必须先 Flush假设 Buffer 里有 60 条数据只存在 TaskManager 内存中。如果 Flink 先把 Source 进度写进 Checkpoint再让 JDBC 慢慢写Checkpoint 已记住 Source 读过这 60 条 JDBC Buffer 却还没有写进 MySQL TaskManager 突然故障内存 Buffer 消失恢复后 Source 可能从这 60 条之后继续MySQL 就永久缺数据。因此顺序必须是先 Flush 当前 JDBC Buffer Flush 成功后当前 Sink 的 snapshotState 才能返回这可以防止 Checkpoint 越过尚未写出的内存数据。四、先 Flush 为什么仍然不是端到端原子事务问题在于 MySQL 和 Flink Checkpoint 没有共用一笔事务。普通 SQL JDBC Sink 做的是executeBatch() - MySQL 已经提交或可见 snapshotState() - 当前子任务继续参与 Checkpoint Checkpoint 全局完成 - 可能还要等其他任务如果 MySQL 已成功但其他算子在全局完成前失败MySQL 不会自动回滚 Flink 也不能使用尚未全局完成的 Checkpoint作业只能回到上一个成功 Checkpoint导致已经写入 MySQL 的范围再次被读取。五、完整看一次 Checkpoint 41 到 42 的故障先设定Checkpoint 41 已经全局成功 Checkpoint 42 尚未全局成功第一次运行T1 从 Checkpoint 41 继续 T2 Paimon Source 发出 Batch A T3 JDBC Sink 将 Batch A Flush 到 MySQL T4 Checkpoint 42 开始 T5 某个 TaskManager 故障Checkpoint 42 失败恢复运行T6 作业只能从 Checkpoint 41 恢复 T7 Source 重新发出 Batch A T8 JDBC Sink 再执行一次 Batch AMySQL 在 T3 已经看到 Batch AT8 又收到相同动作。这不是 Connector 随机重复而是 At-Least-Once 恢复模型的正常结果。六、四个故障位置分别会怎样故障位置MySQL 状态恢复后可能发生什么记录还在 Buffer尚未 FlushMySQL 未看到从成功 Checkpoint 恢复后重新读取并写入executeBatch()执行到一半部分行可能已生效结果不确定整个 Buffer 可能重试或恢复后重放Batch 成功但 Checkpoint 未全局成功MySQL 已看到从旧 Checkpoint 恢复后再次执行Checkpoint 已全局成功MySQL 已看到正常恢复不会回到更早范围除非人为选择旧状态第二和第三种最需要目标端幂等。七、At-Least-Once 到底保证什么At-Least-Once 可以拆成至少送到一次 可能送到多次它重点避免的是数据无声丢失而不是避免重复调用。与几个常见概念对比概念关注点At-Most-Once最多一次可能丢但不重试At-Least-Once至少一次尽量不丢但可能重复Exactly-Once Delivery外部系统观察到每条动作只提交一次Idempotent Result动作可重复执行但最终状态相同普通 JDBC Sink 依赖的是At-Least-Once Delivery Idempotent Result不要把“最终表没有重复行”误写成“每条 SQL Exactly-Once 执行”。八、UPSERT 为什么能让重复执行收敛Batch A 中有U[1001,100.00,PAID]第一次执行INSERTINTOorders_rt(id,amount,status)VALUES(1001,100.00,PAID)ONDUPLICATEKEYUPDATEamountVALUES(amount),statusVALUES(status);MySQL 最终是1001, 100.00, PAID恢复后用相同参数再执行一次结果仍然是1001, 100.00, PAID不会多出第二行也不会把金额再加 100。九、哪些 UPSERT 写法不幂等默认 Dialect 使用赋值amountVALUES(amount)重复执行结果相同。如果自定义成amountamountVALUES(amount)同一条记录重放两次会累加两次。还有一些隐蔽的非幂等行为每次 UPDATE 都让version version 1每次执行都写新的当前时间Trigger 每次插入一条审计记录UPSERT 同时调用外部通知自定义函数每次产生随机值。因此幂等判断必须看完整 DML 和副作用不能只看 SQL 名字里有 UPSERT。十、DELETE 为什么也能结果幂等第一次DELETEFROMorders_rtWHEREid1001;可能影响 1 行。恢复后第二次执行相同 DELETE可能影响 0 行。但最终状态都是orders_rt 中不存在 1001这叫最终结果幂等不代表两次数据库返回的影响行数相同。如果 DELETE Trigger 会记录每次调用就还要单独检查 Trigger 的副作用。十一、sink.max-retries也会造成重复重复不只来自 Checkpoint 恢复。JdbcOutputFormat.flush()会按sink.max-retries重试for(inti0;iexecutionOptions.getMaxRetries();i){try{attemptFlush();batchCount0;break;}catch(SQLExceptione){...}}最简单的失败是连接前就失败 - 一条都没写最麻烦的是MySQL 已执行部分或全部参数 返回结果前网络中断 Flink 不知道哪些已经成功Connector 无法凭空判断数据库真实状态只能重试当前 Buffer。这叫“不确定成功”。UPSERT 和 DELETE 正是为这种情况提供收敛能力。十二、一次 Flush 中还有“前半成功”的可能主键执行器会upsertExecutor.executeBatch();deleteExecutor.executeBatch();假设UPSERT Batch 成功 DELETE Batch 失败Connector 重试或作业恢复后UPSERT 组可能再次执行。所以不能把一次 flush()理解成一笔与 Checkpoint 绑定、全部成功或全部回滚的 MySQL 事务十三、目标主键是幂等成立的第一道门Flink DDL 中PRIMARYKEY(id)NOTENFORCEDMySQL 物理表中也必须有PRIMARYKEY(id)如果 MySQL 没有真实唯一约束第一次 INSERT 一行 恢复后又 INSERT 一行最终出现重复数据。Flink 的NOT ENFORCED不会替 MySQL 检查或创建约束。十四、Key 语义一致是第二道门Paimon、Flink Sink、MySQL 三边 Key 必须一致。例如源端(tenant_id, order_id)目标端却只有order_id恢复重放时即使没有重复行也可能把另一个租户的数据覆盖掉。“没有重复”不等于“数据正确”。十五、多个 Writer 是第三道门UPSERT 对同一批参数重复执行是幂等的但对并发竞争不是自动安全的。Flink 写入 PAID 业务服务写入 REFUNDED Flink 恢复重放 PAID最后可能重新变成 PAID。如果目标表允许其他 Writer 修改同一 Key需要版本号条件更新明确单 Writer事件时间比较冲突检测把出仓结果写入独立服务表。十六、从 Checkpoint 恢复和无状态重跑不是一回事从 Checkpoint / Savepoint 恢复Paimon Source 恢复 Split 和 Snapshot 读取位置。通常只会重放上一个成功 Checkpoint 之后的范围JDBC Sink 依靠主键动作收敛。没有任何状态重新使用latest-fullSource 会重新输出 Paimon 当前仍存在的完整状态。已有 Key 可以通过 UPSERT 覆盖到最新值不会插出重复主键。但有一个重要缺口当前全量只包含 Paimon 仍存在的行不会自动告诉 MySQL 哪些孤儿行应该删除。例如 MySQL 残留id9999Paimon 当前已经没有 9999。重新 latest-full 时没有这条数据也就不会产生对应 DELETE。因此无状态重跑后仍需要主键对账和孤儿清理。十七、Consumer ID 也不能替代 Checkpointconsumer-id有助于管理 Paimon 消费进度和 Snapshot 保留。但 Flink 作业恢复还依赖Enumerator 状态 待分配 Split Reader 中 Split 进度 算子状态 Checkpoint ID没有可用 Checkpoint仅保留相同consumer-id不能自动还原整个作业的精确处理现场。生产设计应该同时考虑consumer-id 稳定 Checkpoint 持续成功 Savepoint 运维流程可用 Snapshot 保留时间覆盖最长恢复窗口十八、普通 SQL JDBC Sink 和 XA Exactly-Once 的区别Flink JDBC DataStream API 还提供基于 XA 的 Exactly-Once Sink 能力。它需要 XADataSource并让数据库事务参与 Checkpoint 的 prepare / commit 协调。但要注意普通 CREATE TABLE ... WITH (connectorjdbc) 并不会因为打开 Checkpoint 就自动变成 XA Exactly-OnceXA 还会引入数据库 XA 支持和配置预提交事务资源占用恢复和事务超时连接池与唯一 XID更复杂的故障处理。如果目标只是维护一张主键服务表At-Least-Once 幂等 UPSERT 往往更简单。如果每次外部副作用都绝不能重复才需要评估 XA、事务 Outbox、去重表或其他端到端方案。十九、怎样设计一个可验证的故障实验不要在生产库直接拔网络。先在测试环境准备Flink Checkpoint 间隔例如 30 秒 JDBC Flush 间隔例如 1 秒 Sink 并行度1 MySQL 测试表有真实主键无 Trigger 测试数据连续递增的 order_id 和 batch_id实验步骤1. 记录最近成功 Checkpoint ID 2. 连续写入一批带 batch_id 的订单 3. 确认 MySQL 已看到这批数据 4. 在下一个 Checkpoint 全局完成前终止当前测试 TaskManager 5. 等作业自动恢复 6. 观察相同 batch_id 是否再次执行但主表最终是否仍正确如果需要证明 SQL 确实重复执行可以使用专门的测试审计机制不要在生产表加非幂等 Trigger 只为观察实验。二十、故障实验至少核对六个结果检查项期望作业恢复点回到最后一次全局成功 CheckpointSource可能重新发出失败窗口内数据MySQL 主键表没有重复 Key值为最新期望状态DELETE重复删除后目标仍不存在Checkpoint恢复后重新持续成功副作用表如果存在需要单独检查是否重复二十一、线上重复、漏数或旧值覆盖怎么排查现象优先检查故障后出现重复行MySQL 是否有真实唯一约束主表不重复审计表重复Trigger 或副作用是否非幂等恢复后旧值覆盖新值是否存在其他 Writer 或版本冲突Checkpoint 一直失败JDBC Flush 延迟、连接错误、MySQL 锁等待无状态重跑后仍有旧行latest-full 不会产生目标孤儿 DELETE少数 Key 删除失败三边主键语义、复合 Key、上游-D二十二、幂等恢复成立的检查清单[ ] MySQL 有真实 PRIMARY KEY / UNIQUE KEY [ ] Paimon、Flink Sink、MySQL 三边 Key 一致 [ ] UPSERT 是覆盖赋值不是累加 [ ] DELETE 只依赖完整主键 [ ] Trigger、审计、通知等副作用可重复执行 [ ] 没有未控制的其他 Writer [ ] Checkpoint 能持续成功 [ ] 无状态重建有孤儿 Key 对账方案二十三、最后记住两句话第一句Checkpoint 前 Flush 是为了防止未写出的内存数据丢失不是为了让 MySQL 和 Checkpoint 自动组成一笔原子事务。第二句普通 JDBC Sink 允许数据重放主键 UPSERT 和 DELETE 解决的是重复执行后的结果收敛不是保证 SQL 从未重复执行。下一篇进入生产调优Batch 设多大 Flush 间隔设多久 并行度为什么不是越高越好 Checkpoint 为什么会形成写入尖峰 MySQL 变慢为什么会让 Flink 背压 怎样做源端与目标端对账本篇关键源码与资料位置GenericJdbcSinkFunction.java在snapshotState()中调用outputFormat.flush()JdbcOutputFormat.javaFlush 重试、连接更新和batchCount清理TableBufferReducedStatementExecutor.javaUPSERT Batch 与 DELETE Batch 的执行顺序MySqlDialect.java幂等赋值式 UPSERT SQLContinuousFileSplitEnumerator.java保存待处理 Split 与下一个 Snapshot 位置FileStoreSourceReader.java保存 Reader 中 Split 消费状态Flink JDBC DataStream 文档普通 Sink 与 XA Exactly-Once Sink 的区别本文基于 Apache Paimon 1.4.2、Apache Flink 1.20.1 和 Flink JDBC Connector 3.3.0-1.20。数据库事务配置、Driver 自动提交与 XA 行为需要结合实际版本和部署方式验证。