
最近把一个实时链路从尽力而为往精确一次上推踩的坑几乎全集中在 Checkpoint 和外部系统的交界处。Flink 的 Checkpoint 本身解决的是作业内部状态的一致性——算子状态、KeyedState、OperatorState 都能靠一次全局快照对齐。但数据一旦要写进 Kafka、MySQL、Doris 这类外部系统内部一致就不等于端到端一致了作业重启后状态回滚了可外部系统里那批数据已经落库重复就发生了。这时候就得靠 Flink CheckPoint 之两阶段提交协议Two-Phase Commit Protocol来兜底——先在外部系统里预占位等 Checkpoint 真正完成后再转正中途挂掉就整批回滚。这篇东西不打算复述官方文档而是把我自己在生产里配过、调过、翻过车的部分完整摊开两阶段提交在 Flink 里的实现骨架长什么样Kafka Sink 内置的事务机制怎么和 Checkpoint 对齐手写一个 MySQL 2PC Sink 的完整代码和参数怎么定以及那些只会出现在日志里的报错到底在说什么。适合已经能跑通 Flink 作业、想把一致性语义从 at-least-once 提到 exactly-once 的同学如果你刚开始学 Flink看到检查点间隔事务超时这些词有点懵也可以先看第一、二节把机制吃透再动手。1. 端到端精确一次到底卡在哪一步1.1 三种一致性语义的真实边界先把概念理清楚不然配参数全是瞎猜。Flink 里说的一致性语义其实分两层很多人混着讲结果调优时找不到北。第一层是作业内部状态一致性。这一层 Flink 自己就能保证Checkpoint 触发时Source 记录当前消费位点算子把自己的状态快照出去所有快照拼成一个全局一致的 Checkpoint。失败重启时从最近一次成功的 Checkpoint 恢复内部状态不会错乱。这一层和你用不用事务完全无关。第二层是端到端一致性也就是从 Source 读进来的数据经过处理写进 Sink整体上恰好处理一次。这一层光靠 Checkpoint 是不够的因为 Sink 写出去的数据在 Flink 状态之外回滚 Checkpoint 不会把已经发出去的消息收回来。于是就有了三种语义语义数据丢失数据重复典型场景at-most-once可能不会日志采集丢几条无所谓at-least-once不会可能大部分实时链路下游能去重exactly-once不会不会计费、对账、库存、指标汇总注意exactly-once 说的是Flink 这套链路内部恰好一次。如果 Sink 是 MySQL而你的业务代码在别处也往同一张表写那整体还是可能重复。别把 exactly-once 当成万能承诺。1.2 为什么单靠幂等写入不够有人会问既然重复会带来问题那我用INSERT ... ON DUPLICATE KEY UPDATE做幂等写入不就行了何必搞两阶段提交这个思路在很多场景下确实够用而且成本低得多。但它有两个前提一是数据必须有天然主键比如订单号、设备 ID 时间戳二是下游必须支持原子 upsert。问题在于很多实时场景的数据是聚合结果比如某商品每分钟的成交额它不是一条可以 upsert 的记录而是一个累加值。这时候重放一次就等于多算一遍幂等就失效了。另一种常见场景是追加型数据比如把处理后的明细写入 Kafka 供下游消费。Kafka 的消息没有主键重复就是实打实的重复消费。两阶段提交解决的正是这两类问题它不依赖下游的幂等能力而是在这个 Checkpoint 到底算不算数这件事上做文章——Checkpoint 成功则这批数据全部可见Checkpoint 失败则这批数据全部不可见。要么全有要么全无这就是原子性。代价也很明确延迟。数据必须先预写等下一次 Checkpoint 完成才能对下游可见。如果你的 Checkpoint 间隔是 1 分钟那么最坏情况下数据要等 1 分多钟才能被下游读到。对延迟敏感的下游这一条就足以否决整个方案。1.3 Checkpoint 存储为什么必须是共享存储顺带说一个被问得最多的问题Flink 是不是一定要 HDFS严格说不是一定但 Checkpoint 的存储必须满足两个条件所有 TaskManager 都能访问且作业重启后依然存在。本地文件系统只在单机伪分布式下勉强能用一旦集群多节点部署TaskManager 各自写本地目录恢复时另一个节点读不到直接报找不到 Checkpoint 元数据。所以生产上通常选 HDFS、对象存储这类共享存储。小规模测试也可以用 NFS 挂载目录或者高度可用的分布式文件系统。关键是别把state.checkpoints.dir配成一个只有本机能访问的路径然后在报错里找半天。2. 两阶段提交在 Flink 里的实现骨架2.1 四个动作begin、preCommit、commit、abortFlink 的两阶段提交抽象在TwoPhaseCommitSinkFunction这个基类里1.15 之后官方推荐迁移到 Sink V2 的Committer接口但思路完全一致老代码现在仍然大量存在。它把一次完整的事务拆成四个动作对应 Checkpoint 生命周期的四个时刻beginTransaction开启一个新事务拿到事务句柄比如数据库连接、Kafka 的 transactionalId。它在算子初始化时和每次commit/abort之后被调用。preCommit在 Checkpoint 触发时调用。把当前事务里攒下的数据预提交同时把事务句柄写进算子状态随 Checkpoint 一起持久化。commit在notifyCheckpointComplete回调里调用也就是 Checkpoint 被 JobManager 确认完成后。这一步才让数据真正对外可见。abortCheckpoint 失败或作业取消时调用丢弃当前事务里未提交的数据。这四个动作的调用顺序就是保证原子性的全部秘密。你可以把它类比成银行转账钱先从 A 账户扣走进入冻结中状态preCommit等对方账户确认能收款了再真正解冻入账commit中途任何一步失败就把冻结的钱退回abort。2.2 Checkpoint 与事务的时序对齐光看方法名还是抽象把时间轴拉出来就清楚了。假设 Checkpoint 间隔 30 秒作业从启动到第三次 Checkpoint时刻Flink 动作Sink 事务状态t0作业启动beginTransaction开启 T1t10s数据持续写入数据写入 T1未提交t30s触发 CP-1preCommit(T1)T1 句柄写入状态t32sCP-1 完成通知到达commit(T1)T1 数据可见同时 beginTransaction 开启 T2t45s数据继续写入数据写入 T2t60s触发 CP-2preCommit(T2)此时若失败则 abort(T2)关键点在于Checkpoint 成功之前事务绝不能 commit。因为 Checkpoint 可能失败需要回滚而一旦 commit 就没法收回了。反过来Checkpoint 成功之后事务必须尽快 commit否则数据一直不可见还会撞上事务超时。还有一个容易忽略的细节preCommit不只是打个标记它必须把事务句柄写进ListState并随 Checkpoint 一起落盘。这样作业挂掉重启后Flink 才能从 Checkpoint 里读出当时有一个事务 T2 处于预提交状态从而决定是补 commit 还是 abort。这个状态里保存未决事务的设计才是两阶段提交能在分布式环境下活下来的原因——它把事务的决策权交给了 Checkpoint 本身。2.3 事务超时和 Checkpoint 间隔必须匹配这是最容易踩的坑没有之一。事务超时Kafka 里是transaction.timeout.ms数据库里通常是连接空闲超时或锁等待超时定义了一个事务多久没动静就自动被判死刑。如果事务超时小于 Checkpoint 间隔就会出现这种尴尬局面事务刚开启没多久就被服务端强制终止了等 Checkpoint 完成去 commit 时报事务已过期或者找不到对应的事务数据直接丢失。安全的下界怎么算transaction.timeout.ms checkpoint.interval 最大单次 Checkpoint 耗时 作业重启耗时余量举个例子Checkpoint 间隔 60 秒大状态下单次 Checkpoint 可能耗时 40 秒恢复一个作业大概 30 秒。那么事务超时至少要大于 130 秒实际我会配到 5 分钟以上也就是 300000 毫秒留出足够余量。为什么留这么多因为在 Checkpoint 失败重试期间当前事务会一直挂着不动。如果配了tolerable-failed-checkpoints 3那可能要连着失败三次才触发重启这段时间事务一直处于开着但没写入的状态。超时值必须能覆盖这个最长悬挂时间。提示反过来事务超时也不能配得过大。事务开太久Kafka 侧会占用事务协调器资源数据库侧会长时间持有锁甚至撑爆max_connections。几分钟到十几分钟是比较务实的区间。3. 手写一个 Kafka 到 MySQL 的端到端精确一次链路3.1 依赖和环境准备光讲原理容易飘直接上手写一遍最清楚。目标链路是Kafka 读取订单消息做简单聚合写入 MySQL要求端到端 exactly-once。依赖上需要dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version${flink.version}/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency环境上Checkpoint 存储要提前配好state.backend: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints state.savepoints.dir: hdfs:///flink/savepoints execution.checkpointing.interval: 60s execution.checkpointing.timeout: 10min execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.max-concurrent-checkpoints: 1 execution.checkpointing.tolerable-failed-checkpoints: 3这里max-concurrent-checkpoints: 1是使用两阶段提交时的硬性建议。允许并发 Checkpoint 意味着可能有多个事务同时处于预提交状态事务之间的提交顺序无法保证下游如果对顺序敏感就会乱。我自己的做法是直接锁死为 1牺牲一点 Checkpoint 吞吐换确定性。3.2 MySQL 两阶段提交 Sink 的实现MySQL 本身没有事务预提交这种语义所以要用连接本身的事务来模拟beginTransaction时setAutoCommit(false)并开启事务preCommit时执行flush()但不 commitcommit时才真正connection.commit()。public class MysqlTwoPhaseCommitSink extends TwoPhaseCommitSinkFunctionOrderStat, Connection, Void { private final String jdbcUrl; private final String user; private final String password; public MysqlTwoPhaseCommitSink(String jdbcUrl, String user, String password) { super(new SimpleVersionedSerializerConnection() { Override public int getVersion() { return 1; } Override public byte[] serialize(Connection c) { // 连接对象不可序列化只保存一个空标记 return new byte[0]; } Override public Connection deserialize(int version, byte[] data) { return null; } }, VoidSerializer.INSTANCE); this.jdbcUrl jdbcUrl; this.user user; this.password password; } Override protected Connection beginTransaction() throws Exception { Connection conn DriverManager.getConnection(jdbcUrl, user, password); conn.setAutoCommit(false); return conn; } Override protected void invoke(Connection conn, OrderStat value, Context context) throws Exception { PreparedStatement ps conn.prepareStatement( INSERT INTO order_stat(order_id, amount, stat_time) VALUES(?,?,?)); ps.setString(1, value.getOrderId()); ps.setBigDecimal(2, value.getAmount()); ps.setLong(3, value.getStatTime()); ps.executeUpdate(); ps.close(); } Override protected void preCommit(Connection conn) throws Exception { // 不做任何提交动作因为数据已经在事务里了 // 这里可以做 flush确保网络缓冲区数据发出 conn.setAutoCommit(false); } Override protected void commit(Connection conn) { try { conn.commit(); } catch (SQLException e) { throw new RuntimeException(commit failed, e); } finally { closeQuietly(conn); } } Override protected void abort(Connection conn) { try { conn.rollback(); } catch (SQLException ignored) { } finally { closeQuietly(conn); } } }这段代码有几个地方值得掰开说。第一Connection不能序列化。TwoPhaseCommitSinkFunction要求事务句柄能进状态但 JDBC 连接是活对象塞不进去。所以序列化器里只写一个空字节数组靠 Flink 在恢复时重新建立连接。这也意味着恢复后的事务语义是重新执行一遍未提交的数据而不是续上原连接。这一点必须接受否则整个模型不成立。第二preCommit里几乎什么都不用做。因为 MySQL 的事务本身就有未提交不可见的特性天然满足第一阶段的要求。这和 Kafka 不一样Kafka 需要显式调flush()把缓冲消息发出去让 broker 端持有但不标记为已提交。第三commit抛异常会导致作业失败并重启。这是有意的commit 阶段失败不能静默吞掉否则数据就丢了。让作业失败、从上一个 Checkpoint 恢复、重新走一遍流程才是正确姿势。3.3 主程序与参数配置把 Sink 接进作业StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000); env.getCheckpointConfig().setCheckpointTimeout(600_000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3); KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka:9092) .setTopics(order-topic) .setGroupId(order-stat-group) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamOrderStat stats env .fromSource(source, WatermarkStrategy.noWatermarks(), kafka-source) .map(new ParseAndAggregateFunction()) .name(aggregate) .uid(aggregate-uid); stats.addSink(new MysqlTwoPhaseCommitSink( jdbc:mysql://mysql:3306/dw?useSSLfalse, flink, flink_pwd)) .name(mysql-2pc-sink) .uid(mysql-2pc-sink-uid); env.execute(order-stat-exactly-once);两个参数特别想强调。minPauseBetweenCheckpoints设成 Checkpoint 间隔的一半是为了给 commit 操作留出窗口期。两阶段提交的 commit 是发生在 Checkpoint 完成回调里的如果 Checkpoint 一个接一个连轴转commit 还没执行完下一个 Checkpoint 就来了事务会堆积。uid必须显式指定而且上线后绝对不能改。Flink 靠 uid 把算子和状态做映射uid 变了就相当于换了个算子之前保存的事务句柄状态全部丢失。后果是那些处于预提交状态的事务既不会被 commit 也不会被 abort永久悬挂在数据库里占着锁。3.4 怎么验证真的做到了精确一次写完不算完得能证明。我的验证方法是故意制造故障 对数。第一步正常跑 5 分钟记录 MySQL 里的总行数记为 A。第二步重启作业让它重新消费一部分已经处理过的数据把 Kafka 消费位点人为往前调一点继续跑 5 分钟。第三步等作业稳定后再次统计行数记为 B。如果 B 和 A 的差值恰好等于新增数据量说明没有重复如果 B 明显偏大说明有两阶段提交没生效的地方如果 B 偏小那就更严重了是丢数。另一个更直接的办法是打开网络抓包或者在 Sink 的commit/abort里打点统计两个方法被调用的次数。正常情况下commit次数应该等于成功的 Checkpoint 次数可能少一次因为最后一次 Checkpoint 未必完成abort次数应该只在故障时出现。注意做这类验证一定要在独立的测试环境做。调整消费位点会污染生产数据代价可能是几小时的对账工作。4. Kafka Sink 内置的两阶段提交与实战陷阱4.1 Kafka 自己的事务机制怎么和 Flink 接上Kafka 从 0.11 版本开始支持事务Flink 的 Kafka Sink 直接复用了它。开启方式很直接KafkaSinkString sink KafkaSink.Stringbuilder() .setBootstrapServers(kafka:9092) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(result-topic) .setValueSerializationSchema(new SimpleStringSchema()) .build()) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(order-stat-) .setProperty(transaction.timeout.ms, 300000) .build();底层发生了什么每个并行 Sink 子任务会用transactionalIdPrefix subtaskIndex拼出一个唯一的 transactionalId向事务协调器注册。Checkpoint 触发时调producer.flush()Agent 把缓冲的消息发给 broker但消息处于未提交状态消费者用read_committed隔离级别看不到。notifyCheckpointComplete到达后调producer.commitTransaction()这批消息才对外可见。transactionalIdPrefix有两个约束同一时刻不能有两条作业用同一个前缀否则第二个作业启动时会因为 transactionalId 冲突而失败上线后不能随便改改了等于换了事务身份之前未提交的事务就成孤儿了。4.2 那些让人抓狂的报错报错一InvalidProducerEpochException或者ProducerFencedException这个报错的含义是你的事务身份被别人抢了。最常见的原因是同一个作业被重复提交了两次两个实例用同样的 transactionalId 在跑。排查方向检查调度平台上是否有残留的僵尸作业或者作业重启时旧实例还没完全退出。报错二InvalidTxnStateException或者commitTransaction超时一般是事务超时了。要么是transaction.timeout.ms配得太小要么是 Checkpoint 长时间卡住导致事务悬挂超过阈值。我遇到过一次正是下游 HDFS 集群抖动Checkpoint 卡了 8 分钟直接超了默认的 5 分钟事务超时。解决办法有两个调大超时值或者缩短 Checkpoint 超时时间让它快速失败重试。报错三下游一直读不到数据事务提交了但下游看不到八成是消费者用的隔离级别不对。默认的read_uncommitted或不设置隔离级别时行为不确定正确做法是显式配置props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, read_committed);用read_committed之后消费者只能看到已提交事务的消息未提交的会被过滤掉。代价是会有一定的读取延迟。报错四TimeoutException: Expiring N record(s)消息在缓冲区里等太久被丢弃了。检查linger.ms和batch.size的配合也检查事务里是否积累了过多数据导致发送耗时超过阈值。两阶段提交下一个事务里塞的数据量不宜过大——事务跨越时间越长占用的缓冲区、锁和协调器资源越多。5. 常见问题与排查技巧实录5.1 问题速查表现象最可能原因排查入口处理方式作业重启后数据重复Sink 未开启两阶段提交检查 Sink 是否实现CheckpointListener换用支持 2PC 的 Sink 或下游加幂等下游长时间读不到数据事务已开启但 Checkpoint 未完成看 Checkpoint 成功率与耗时排查 Checkpoint 卡点缩短间隔报事务超时事务超时小于 Checkpoint 周期对比transaction.timeout.ms与间隔超时值放大到 3 到 5 倍间隔事务永久悬挂uid 被改动或作业异常终止查外部系统未决事务列表手动终止悬挂事务固定 uid连接数暴涨abort 未正确释放连接统计show processlist在 finally 块里兜底关闭Checkpoint 长期失败状态过大或反压严重Checkpoint 详情页的算子耗时拆大状态、开增量 Checkpoint提交后仍有重复下游同时有别的写入源全链路梳理写入口统一收口到 Flink 链路5.2 几条只有踩过才知道的经验第一事务超时宁可配大不配小。配小了丢数据的代价是几小时的对账配大了顶多多占点资源。这条经验背后是一次真实的教训我们把事务超时配成了 2 分钟Checkpoint 间隔 1 分钟理论上是够的但一次机房网络抖动让单次 Checkpoint 花了 3 分半结果那一批数据在 commit 时直接报事务不存在。损失不大但定位花了整整一个下午。第二Sink 的并行度不要盲目调大。每个并行子任务对应一个独立事务、一个独立连接。并行度调到 16就意味着最多有 16 个事务同时挂在数据库上连接池配置跟不上就是一堆Too many connections。除非下游确实扛不住写入压力否则 Sink 并行度保持在 2 到 4 是比较稳的。第三abort 一定要做资源兜底。很多人写abort时只调了rollback()忘了关连接。作业频繁重启时这些没关掉的连接会累积最终把数据库连接数吃满。正确写法是try { rollback } catch {} finally { close }关连接这一步不能省。第四不要在 Sink 里做重业务逻辑。我见过有人在invoke里调用外部 HTTP 接口做数据补全一旦接口超时整个事务就卡住了Checkpoint 跟着失败。两阶段提交的窗口期很宝贵invoke里只应该做纯粹的写入动作任何可能阻塞的操作都要挪到上游算子。第五测试环境一定要模拟故障。生产上第一次遇到 commit 阶段崩溃时如果没演练过心态很容易崩。我的做法是在测试环境用kill -9直接杀 TaskManager观察作业重启后数据库里的事务是被正确提交还是回滚是否有残留。跑通几次之后对这套机制的行为就心里有数了。6. 选型什么时候该上两阶段提交6.1 幂等写入和两阶段提交的取舍不是所有场景都值得上两阶段提交。判断标准其实很简单问自己两个问题数据有没有天然主键下游能不能接受 1 到 2 个 Checkpoint 间隔的可见延迟如果数据有天然主键、下游支持 upsert那用幂等写入就够了简单、延迟低、还不用管事务超时这一堆麻烦。比如订单表的同步订单号本身就是主键直接INSERT ... ON DUPLICATE KEY UPDATE完事。如果数据是聚合值、是追加型消息、或者下游明确要求不能看到中间态那就得上两阶段提交。指标类、对账类、计费类业务基本都属于这一类。还有一类是混合方案值得单独提一句Source 端保证精确一次Sink 端用幂等。Kafka Source 本身能通过位点提交保证精确一次Sink 端如果下游有主键就做幂等。这种组合在很多场景下已经够了比全链路 2PC 简单得多。6.2 各类连接器的现状与坑点实际选型时最大的痛点是不同连接器对两阶段提交的支持程度差异极大而且这个差异在文档里往往一句话带过。Kafka 连接器支持最好内置DeliveryGuarantee.EXACTLY_ONCE开箱可用也是我推荐的入门练习对象。JDBC 连接器的情况要复杂一些。官方 JDBC Sink 在较早版本里对 exactly-once 的支持是通过TwoPhaseCommitSinkFunction实现的但需要数据库开启 XA 支持比如 MySQL 的 XA 事务。而 XA 事务在生产里有很多争议性能损耗明显、长时间悬挂的 XA 事务需要 DBA 手动清理、某些云数据库甚至不开放 XA 权限。我自己的经验是除非业务强需求否则 JDBC 链路优先用幂等 主键的方案而不是硬上 2PC。如果用 JDBC 遇到了连接层面的异常比如驱动版本与数据库版本不匹配导致的元数据读取错误先确认驱动版本再确认连接串参数最后才怀疑两阶段提交的配置。Doris 连接器这几年的成熟度提升很快。它的 Sink 采用 Stream Load 方式写入本身通过 Label 机制做幂等——同一个 Label 的导入请求重复提交会被 Doris 自动去重。所以它的 exactly-once 实现思路和两阶段提交不完全一样更偏向用 Label 做幂等。需要注意的是用 Flink CDC 写入 Doris 时经常遇到类型映射问题比如源端是日期类型Doris 侧字段类型不匹配报出类似类型不一致的元数据错误。这类问题本质上是类型系统对不上和事务机制无关但排查时容易和一致性配置混在一起建议先单独把类型对齐再验证一致性。TiDB 通过 Flink SQL 写入的链路因为有分布式事务支持理论上是比较容易做精确一次的。用 Flink SQL 的话sink表配置里开启相关语义就行但要注意 TiDB 侧事务大小限制一批写入量过大时会报事务过大失败。提示选连接器之前先去对应版本的官方文档确认它到底支持哪种语义。很多数据重复的锅最后都不是两阶段提交的配置问题而是连接器压根就没提供 exactly-once 能力。6.3 一个务实的落地顺序如果现在就要把这个方案推上线我建议按这个顺序来别一上来就全链路开 2PC。先做一轮链路梳理把所有写入口列出来确认哪些是主链路、哪些是旁路。然后从主链路里挑一条数据量适中的先开两阶段提交用小流量跑一周观察 Checkpoint 成功率、commit 耗时、下游可见延迟这三个指标。指标稳定了再逐步扩大范围。同时把监控补上。要盯的东西包括Checkpoint 成功率与平均耗时、Sink 的 commit 与 abort 调用次数可选、外部系统里的未决事务数量、数据库连接数。这几个指标里未决事务数量是最灵敏的预警信号一旦持续增长就说明 commit 环节出了问题得赶紧查。最后再分享一个小技巧。两阶段提交最难排查的情况是事务悬挂因为它不报错、不告警只是数据凭空少了。我的做法是在 Sink 里给每个事务起一个可读的标识比如jobName-checkpointId-subtaskIndex并把这个标识写进外部系统的备注字段或者日志。出问题时直接拿这个标识去外部系统里搜能立刻定位到是哪个 Checkpoint、哪个子任务的事务没被处理比翻 Flink 日志快得多。这个标识的生成逻辑要放在beginTransaction里并且随 Checkpoint 一起持久化保证重启后依然能对上。