在广告实时计费里,一条曝光被处理两次,就是账上多算了一笔钱——多计费轻则对账对不平、财务月底追着你要说法,重则客户投诉超扣、平台被迫赔付。所以”作业重启后数据别重复”这件事,在计费/结算链路上不是”锦上添花”,而是红线。而这条红线的技术名字,就叫 exactly-once(精确一次)。
问题是,“精确一次”这四个字被用得太随意,很多人以为它是”每条消息只被处理一次”——这是个天大的误解,也是本篇要拆开的第一件事。
一句话定位:exactly-once 不是”每条消息只被物理处理一次”,而是”故障重放后,反映到状态与对外结果上等效于一次”(effectively-once);端到端做到它 = 可重放的 source + 一致性快照 + 事务/幂等的 sink,三者缺一不可。
TL;DR
- 三种语义:
at-most-once(可能丢)、at-least-once(可能重复)、exactly-once(状态一致)。计费链路里,丢 = 少收钱、重复 = 多扣钱,两者都要命,所以必须 exactly-once。 - “精确一次”是被误解最深的术语:它指的是 effectively-once——状态与对外结果等效于处理了一次,不是每条消息物理上只流经算子一次。故障时消息照样重放、算子照样重算。
- Flink 内部 exactly-once 靠 checkpoint:基于 Chandy-Lamport 的异步屏障快照做一致性快照,失败时整体回滚到上一个 checkpoint 的状态,并把 source 的读取位点也一起回退、重放(这套一致性快照机制,状态管理与 Checkpoint 容错 会专门深挖)。
- 内部一致 ≠ 端到端一致:Flink 能回滚自己的状态,但回滚不了已经写出去的副作用。要端到端,还需 source 能按 offset 重放 + sink 不产生重复副作用。
- 两条 sink 路线:幂等 sink(写入天然幂等,如按主键 upsert、按
requestId去重)vs 事务 sink(两阶段提交,把”写”和”提交”分开)。 - TwoPhaseCommitSinkFunction:
beginTransaction → preCommit(checkpoint 快照时)→ commit(checkpoint 完成通知时)→ abort(失败回滚),四个动作严格绑定 checkpoint 生命周期。 - Kafka 端到端 exactly-once:
KafkaSource(offset 存进 checkpoint)+KafkaSink设EXACTLY_ONCE(事务性 producer、transactional.id),下游消费者必须isolation.level=read_committed。 - 最坑的一个参数:
transaction.timeout.ms必须 > checkpoint interval + 作业恢复时间,否则事务在 commit 前被 broker 判超时中止,数据直接丢。默认 1 小时上限还受 broker 端transaction.max.timeout.ms约束。 - 代价是延迟:事务 sink 的数据”要等 checkpoint 完成才对下游可见”,端到端延迟 ≈ checkpoint 间隔量级。要低延迟就得缩短 checkpoint 间隔或改用幂等 sink。
- 幂等兜底:拿不到端到端事务时,用
requestId去重(去重表 / 布隆过滤器前置),把”至少一次”收敛成”效果上一次”,见 幂等与一致性。
Table of contents
Open Table of contents
1. 先把三种投递语义说清楚
任何”可靠性”讨论都得先定义”可靠到什么程度”。流处理里有三个层级,区别只在故障重放时会发生什么:
- at-most-once(至多一次):数据被处理 0 或 1 次。不重试、不重放,崩了就丢。延迟最低、实现最简单,但会丢数据。
- at-least-once(至少一次):数据被处理 1 到 N 次。保证不丢,但故障恢复时会重放已经处理过的数据,导致重复。
- exactly-once(精确一次):数据对状态和结果的影响恰好一次。既不丢也不重复——代价是最复杂、延迟也最高。
三种语义的分水岭都在”故障重放时会怎样”:at-most-once 丢、at-least-once 重复、exactly-once 靠快照回滚 + 事务/幂等做到”效果上一次”。
落到广告实时计费上,为什么必须是最贵的那一档?因为另外两档的代价都是真金白银:
- 用 at-most-once:作业崩一次,那段时间的曝光没计费 → 平台少收钱、账单少一截。
- 用 at-least-once:作业重启回放,部分曝光被重复计费 → 要么多扣客户预算(投诉、退款),要么内部对账对不平(财务事故)。
计费是”重复即金钱损失、丢失即收入损失”的典型场景,两个方向都不能容忍——这就是为什么这条链路值得付出 exactly-once 的全部复杂度。
2. 澄清”精确一次”:它说的是状态,不是物理次数
这是全篇最重要的一个观念,先钉死它,后面才不会绕晕。
“exactly-once”并不意味着每条消息只被物理地处理一次。 恰恰相反:一旦发生故障,Flink 会把作业回滚到上一个一致性快照,然后从那个位点重放 source——这意味着快照之后、故障之前处理过的那批消息,会被再处理一遍。算子的 map、process、聚合逻辑,统统会对这批消息重复执行。
那”精确一次”到底精确在哪?精确在最终状态与对外产生的结果上。业界更准确的叫法是 effectively-once(效果上一次):
不管中间重放了多少次、算子重算了多少遍,反映到 Flink 的托管状态和外部 sink 上的净效果,等同于每条消息恰好被处理了一次。
理解这一点,才能理解后面所有机制的设计动机:
- 状态为什么能”精确一次”:因为回滚会把状态一并退回到快照点,重放产生的状态变更是”从快照点重新累加”,而不是”在已有状态上再加一遍”——所以不会重复累计。
- 对外副作用为什么难:因为你已经写到外部系统的那些数据,回滚不掉。重放会让算子再写一遍,如果 sink 不做特殊处理,外部就有了重复。这正是第 3 节要展开的”端到端为什么难”。
一句话记住:Flink 的 exactly-once = 内部状态一致(免费,靠 checkpoint)+ 对外副作用一致(收费,靠幂等或事务 sink)。
3. Flink 内部的 exactly-once:一致性快照 + 回滚 + 重放
内部这一半的完整机制,状态管理与 Checkpoint 容错 会专门深挖,这里只快速串起与 exactly-once 直接相关的三个环节。
3.1 一致性快照:异步屏障快照(ABS)
Flink 的 checkpoint 基于 Chandy-Lamport 分布式快照算法的变体——异步屏障快照(Asynchronous Barrier Snapshotting)。核心是往数据流里周期性注入一种特殊记录 checkpoint barrier:
- barrier 从 source 注入,随数据一起向下游流动,把数据流切成”barrier 之前”和”barrier 之后”两段。
- 每个算子收到 barrier n 时,就对自己当前的状态做一份快照(异步写到状态后端,如 RocksDB / 远端存储)。
- 多输入算子要做 barrier 对齐(alignment):等所有输入通道的 barrier n 都到齐,才快照——这保证快照反映的是”恰好处理完 barrier 之前所有数据”的那个一致性切面。
所有算子都完成快照 n 后,checkpoint n 才全局完成。这份全局快照的语义是:“如果从这里恢复,就像所有算子都恰好处理到 barrier n 之前的最后一条数据”。
3.2 失败回滚 + source 重放:一致性从何而来
当某个 task 失败,Flink 的恢复动作是整体性的(默认 region/full failover):
- 回滚状态:把所有算子的状态恢复到最近一次成功的 checkpoint n。
- 回退 source 位点:checkpoint n 里存了 source 的读取位点(如 Kafka 的 offset)。恢复时 source 从这个 offset 重新开始读。
- 重放:offset n 之后的数据被重新拉取、重新处理。
关键就在第 2、3 步:source 的位点是 checkpoint 状态的一部分。状态回退到 n、offset 也回退到 n,两者是同一份快照里的原子对。于是”状态”和”已消费的数据”永远对得上——重放不会导致状态重复累加,因为状态本身也退回了重放的起点。
这就是内部 exactly-once 的全部秘密:把 source 位点和算子状态绑进同一次一致性快照,失败时一起回滚、一起重放。 它对 source 有一个硬要求——source 必须可重放(replayable),即能按记录下来的位点重新读取历史数据。Kafka(按 offset)、文件系统(按 offset/文件位置)天然满足;而一个”读完即焚”的 socket source 就做不到,这类 source 无法支撑 exactly-once。
4. 端到端为什么难:内部一致 ≠ 端到端一致
到这里,Flink 作业内部的 exactly-once 已经齐活了。但真实作业不是自娱自乐——它要把结果写到外部:写进另一个 Kafka topic、写进数据库、写进计费账本。这一写,麻烦就来了。
设想计费作业:从 Kafka 读曝光事件 → 累加每个广告主的花费 → 把结果写到下游。现在 checkpoint n 之后、checkpoint n+1 之前,作业崩了:
- Flink 会把状态回滚到 checkpoint n,source offset 也退回 n,重放 n 之后的曝光。内部状态没问题。
- 但是:n 之后、崩溃之前,作业已经把一部分计费结果写到下游了。这些写出去的数据,Flink 回滚不掉——它们已经在外部系统里,已经被别人读走了。
- 重放时,这批曝光被再算一遍、再写一遍。下游于是收到了重复的计费。
这就是端到端的核心难点:Flink 能回滚自己的状态,却回滚不了已经产生的外部副作用。 内部的一致性快照只解决了”我自己的账对不对”,解决不了”我已经泼出去的水”。
要补上这一半,只有两条路:
- 让 source 可重放(第 3 节已解决)——保证重放的数据源头是确定、可回退的。
- 让 sink 不产生重复副作用——要么写入本身幂等(重复写 = 写一次),要么把”写”延迟到”确定不会回滚”之后(事务)。
这两条路里第 2 条是全部难点所在,下面两节分别讲幂等 sink 和事务 sink。
5. 两种 sink 路线:幂等 vs 事务
5.1 幂等 sink:让”重复写”等于”写一次”
如果外部写入天然幂等——同样的数据写 N 次,效果和写 1 次一样——那重放导致的重复写就无害了。典型手段:
- 按主键 upsert:以业务主键(如
campaignId + 分钟窗口)做INSERT ... ON DUPLICATE KEY UPDATE或REPLACE。重放写同一个 key,只是把值覆盖成同样的结果,不会累加出重复。 - 按去重键写入:给每条曝光带一个全局唯一的
requestId,下游用它做主键或去重表。重复的requestId直接被丢弃。
// 幂等写:按 (campaignId, minuteBucket) upsert 累计值
// 重放时写的是"从快照点重算出的同一个累计值",覆盖即可,不会重复累加
void writeSpend(String campaignId, long minuteBucket, long spendMicros) {
jdbc.update(
"INSERT INTO ad_spend(campaign_id, minute_bucket, spend_micros) " +
"VALUES (?,?,?) ON DUPLICATE KEY UPDATE spend_micros = VALUES(spend_micros)",
campaignId, minuteBucket, spendMicros);
}
注意一个陷阱:幂等的前提是”写的是最终值”而不是”写的是增量”。上面写的是”这个窗口的累计花费”(覆盖语义),幂等成立;如果写的是 UPDATE ... SET spend = spend + ?(增量累加),那重放就会多加一遍,幂等不成立。计费场景尤其要警惕这一点——能用”覆盖最终值”就别用”累加增量”。
幂等 sink 的好处是简单、低延迟(不用等 checkpoint,写完就可见);限制是不是所有 sink 都能做成幂等(比如往 Kafka 追加消息,天然不幂等)。
5.2 事务 sink:把”写”和”提交”分开
当 sink 无法做成幂等(最典型的就是”往下游 Kafka 追加消息”),就得上事务。核心思路:
数据先写进一个”未提交”的事务里(对下游不可见),等到对应的 checkpoint 全局完成、确认永远不会回滚了,再一次性 commit 让它可见;如果中途失败,就 abort 掉整个事务,那些写入被丢弃。
这样,只有”已提交 checkpoint 对应的数据”才对外可见,回滚丢弃的都是未提交的写入——外部永远看不到会被回滚的数据,重复也就无从谈起。这套协议就是两阶段提交(Two-Phase Commit, 2PC),Flink 用 TwoPhaseCommitSinkFunction 抽象了它。
6. TwoPhaseCommitSinkFunction:两阶段提交协议
Flink 把事务 sink 的通用骨架抽成了 TwoPhaseCommitSinkFunction(KafkaSink 的 EXACTLY_ONCE 模式内部就是这套逻辑)。它要求实现四个方法,每个都严格绑定 checkpoint 生命周期的一个时点:
| 方法 | 触发时机 | 干什么 |
|---|---|---|
beginTransaction() | 上一个事务提交后 / 作业启动 | 开启一个新事务,后续处理结果都写进它 |
preCommit() | checkpoint 快照时(snapshotState) | flush 数据、把事务句柄写进 Flink 状态;进入”预提交”,但还不 commit |
commit() | checkpoint 全局完成通知时(notifyCheckpointComplete) | 真正 commit 事务,数据对下游可见 |
abort() | 失败 / 恢复时 | 回滚事务,丢弃未提交的写入 |
把它和 checkpoint 生命周期对齐,时序是这样的:
两阶段提交的”提交点”与 checkpoint 完全对齐:checkpoint 快照时 preCommit(第一阶段),checkpoint 全局完成通知时才 commit(第二阶段)。这保证”只有不会被回滚的数据才对外可见”。
理解这套协议的关键,是看清 commit 为什么必须发生在 notifyCheckpointComplete:
- checkpoint 全局完成 = 这次快照已经持久化、绝对不会再回滚到它之前。这时候 commit,等于宣布”这批数据是永久的、可以让别人看了”。
- 如果在
preCommit(快照时)就 commit,而这个 checkpoint 还没最终成功(可能其他算子快照失败导致整个 checkpoint 作废),那就 commit 了一批”本该被回滚”的数据 → 破坏 exactly-once。 - 所以顺序铁律是:先确认 checkpoint 永久成功,再 commit 事务。 这也是”两阶段”的本质——preCommit 是”投票 yes、准备好但不生效”,commit 是”协调者拍板后才真正生效”。
还有个容错细节:恢复时对”已 preCommit 但还没收到 commit 通知”的事务,要重新 commit(而不是 abort)。因为它对应的 checkpoint 已经成功了(否则不会恢复到它之后),这批数据是该生效的。TwoPhaseCommitSinkFunction 把待提交事务的句柄存进了 Flink 状态,恢复时能拿回来重新 commit——这就是为什么 preCommit 要”把事务句柄写进快照”。
7. Kafka 端到端 exactly-once:实战
把上面的理论落到最常见的 Kafka → Flink → Kafka 链路。这也是 Kafka 可靠性与 exactly-once 在 Flink 侧的落地。
端到端三要素:可重放 source(offset 进 checkpoint)+ 事务 sink(两阶段提交对齐 checkpoint)+ 下游 read_committed。任何一环缺失,端到端 exactly-once 就破。
7.1 KafkaSource:offset 存进 checkpoint
新版 KafkaSource(flink-connector-kafka)默认就把 offset 作为算子状态存进 checkpoint,恢复时从 checkpoint 里的 offset 重放——不要依赖 Kafka 自身的 enable.auto.commit,那套自动提交与 Flink 的 checkpoint 语义不对齐。Flink 只在 checkpoint 完成后把 offset 回写 Kafka(仅供监控/展示),真正用于恢复的是 checkpoint 里那份。
KafkaSource<Impression> source = KafkaSource.<Impression>builder()
.setBootstrapServers("kafka:9092")
.setTopics("ad-impressions")
.setGroupId("billing-job")
.setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
.setValueOnlyDeserializer(new ImpressionDeserializer())
.build();
7.2 KafkaSink:EXACTLY_ONCE + 事务性 producer
KafkaSink 设 DeliveryGuarantee.EXACTLY_ONCE,底层就走 TwoPhaseCommitSinkFunction 那套:用 Kafka 的事务性 producer,每个 checkpoint 一个事务,preCommit 时 flush、commit 时提交事务。必须设 transactionalIdPrefix——Flink 会基于它给每个并行 sink 子任务生成唯一的 transactional.id(Kafka 靠它做事务的 fencing 与恢复)。
KafkaSink<SpendRecord> sink = KafkaSink.<SpendRecord>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("ad-spend")
.setValueSerializationSchema(new SpendSerializer())
.build())
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) // 事务性 sink,两阶段提交
.setTransactionalIdPrefix("ad-spend-sink") // 必填:事务 id 前缀
// 事务超时必须 > checkpoint interval + 恢复时间(见 §8)
.setProperty("transaction.timeout.ms", "900000") // 15 分钟
.build();
env.enableCheckpointing(60_000); // checkpoint 间隔 60s
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.fromSource(source, WatermarkStrategy.noWatermarks(), "kafka-src")
.keyBy(Impression::campaignId)
.process(new SpendAggregator()) // 有状态累加
.sinkTo(sink);
7.3 下游必须 read_committed,否则前功尽弃
这是最容易被漏掉的一环:上游 Flink 用了事务 sink,下游消费者却没设 read_committed,那 exactly-once 直接破功。 因为默认 isolation.level=read_uncommitted 的消费者会读到未提交(甚至将被 abort)的事务数据,重复照旧出现。
# 下游消费 ad-spend 的消费者,必须这样配
isolation.level=read_committed
read_committed 让消费者只读已 commit 的事务消息,未提交/已中止的对它完全不可见。端到端 exactly-once 是一条链上所有环节的共同契约,任何一端不守约都会破。
8. 取舍与坑:延迟、超时、幂等 vs 事务
exactly-once 不是免费午餐,落地前必须认清这几个代价与坑。
8.1 事务 sink 增加端到端延迟
事务 sink 的数据要等对应 checkpoint 全局完成、commit 之后才对下游可见。也就是说,端到端延迟被 checkpoint 间隔”量化”了——checkpoint 间隔 60s,数据最坏要等近 60s 才可见。这是 exactly-once 与低延迟的根本矛盾:
- 要更低延迟 → 缩短 checkpoint 间隔(但太频繁会增加快照开销、加重 反压)。
- 能接受量化延迟 → 保持较长间隔,吞吐更好。
- 真的要亚秒级可见 → 放弃事务 sink,改用幂等 sink(写完即可见,靠幂等消化重复)。
8.2 transaction.timeout.ms:最容易踩的雷
Kafka 的事务有超时时间 transaction.timeout.ms:一个事务开启后,超过这个时间还没 commit,broker 会主动把它中止(abort)。
在 Flink 事务 sink 里,一个事务的存活时间是”从 beginTransaction 到 commit”,而 commit 发生在下一个 checkpoint 完成之后。如果作业发生故障、恢复耗时较长,事务可能从开启到恢复后重新 commit 之间跨了很久。一旦这段时间超过 transaction.timeout.ms,事务被 broker abort,恢复时想重新 commit 却发现事务已经没了 → 那段数据永久丢失。
铁律:
transaction.timeout.ms必须 > checkpoint interval + 作业最坏恢复时间,并留足余量。
同时注意两个约束:Kafka broker 端的 transaction.max.timeout.ms(默认 15 分钟)是上限,producer 端设的值不能超过它——要更大得先调 broker。生产上常见配法是把 checkpoint 间隔设为分钟级、transaction.timeout.ms 设为远大于它(如间隔 1 分钟、超时 15 分钟),给恢复留出充足窗口。
8.3 幂等 vs 事务:怎么选,怎么兜底
| 维度 | 幂等 sink | 事务 sink(2PC) |
|---|---|---|
| 前提 | 写入天然幂等(upsert / 去重键) | sink 支持事务(Kafka / 支持 XA 的 DB) |
| 延迟 | 低,写完即可见 | 高,等 checkpoint 完成才可见 |
| 复杂度 | 低 | 高(事务生命周期、超时、恢复) |
| 典型场景 | 写 DB(按主键覆盖)、写去重表 | 写 Kafka、要求严格事务边界的下游 |
选型经验:能做成幂等就优先幂等——简单、低延迟。只有当 sink 无法幂等(如往 Kafka 追加)、或业务要求严格的事务原子性时,才上事务 sink。
而当两者都拿不到”完美的端到端事务”时,还有一条务实的幂等兜底路线:给每条业务事件带全局唯一 requestId,下游按它去重。
- 去重表:下游维护一张
processed_request_id表,写入前先查、已存在就跳过。可靠但每条都要查库。 - 布隆过滤器前置:用布隆过滤器做第一道去重——“一定不存在”就直接放行、“可能存在”再查去重表确认。布隆过滤器能挡掉绝大多数的”确实是新请求”的查库,把去重表的压力降下来,是高吞吐去重的经典兜底组合。
这套幂等去重的更一般讨论(幂等键设计、去重窗口、最终一致),见 幂等与一致性。
参考
- Apache Flink. Fault Tolerance via State Snapshots(Checkpointing 与 exactly-once 概念):Flink 官方对一致性快照与 exactly-once 语义的说明。
- Apache Flink. Fault Tolerance Guarantees of Data Sources and Sinks(各连接器的语义保证表):官方汇总各 source/sink 能提供的投递语义。
- Apache Flink. Kafka Connector(KafkaSource / KafkaSink 与 EXACTLY_ONCE 配置):Kafka 连接器的 offset、事务、
transactionalIdPrefix与超时配置。 - Piotr Nowojski et al. An Overview of End-to-End Exactly-Once Processing in Apache Flink(TwoPhaseCommitSinkFunction 经典博文):Flink 官方博客对两阶段提交 sink 的权威讲解。
- Chandy, K. M., Lamport, L. Distributed Snapshots: Determining Global States of Distributed Systems:Flink 一致性快照算法的理论源头。