在广告实时计费里,一条曝光被处理两次,账上就多算了一笔钱——轻则对账对不平、财务月底追着你要说法,重则客户投诉超扣、平台被迫赔付。所以「作业重启后数据别重复」这件事,在计费与结算链路上从来不是锦上添花的加分项,而是一条红线,而这条红线的技术名字就叫 exactly-once(精确一次)。
麻烦的是,「精确一次」这四个字被用得太随意,很多人默认它等于「每条消息只被处理一次」。这是个天大的误解,也是本篇要拆开的第一件事。
一句话定位:exactly-once 不是「每条消息只被物理处理一次」,而是「故障重放后,反映到状态与对外结果上等效于一次」(effectively-once);端到端做到它 = 可重放的 source + 一致性快照 + 事务或幂等的 sink,三者缺一不可。
TL;DR
- 三种语义:
at-most-once(可能丢)、at-least-once(可能重复)、exactly-once(状态一致)。计费链路里丢 = 少收钱、重复 = 多扣钱,两个方向都要命,所以只能选最贵的那一档。 - 「精确一次」是被误解最深的术语:它指的是 effectively-once,即状态与对外结果等效于处理了一次,不是每条消息物理上只流经算子一次——故障时消息照样重放、算子照样重算。
- Flink 内部 exactly-once 靠 checkpoint:基于 Chandy-Lamport 的异步屏障快照做一致性快照,失败时整体回滚到上一个 checkpoint 的状态,并把 source 的读取位点一起回退、重放(这套一致性快照机制,状态管理与 Checkpoint 容错 会专门深挖)。
- 内部一致 ≠ 端到端一致:Flink 能回滚自己的状态,却回滚不了已经写出去的副作用,所以还需要 source 能按 offset 重放、sink 不产生重复副作用。
- 两条 sink 路线:幂等 sink(写入天然幂等,如按主键 upsert、按
requestId去重)与事务 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 判超时中止,数据直接丢;Flink 侧默认的 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 上的净效果,等同于每条消息恰好被处理了一次。
把这句话理解透,后面所有机制的设计动机就都顺了:
- 状态为什么能做到精确一次:因为回滚会把状态一并退回到快照点,重放产生的状态变更是「从快照点重新累加」而不是「在已有状态上再加一遍」,所以不会重复累计。
- 对外副作用为什么难:因为已经写进外部系统的那些数据,Flink 回滚不掉。重放会让算子再写一遍,sink 若不做特殊处理,外部就凭空多出一份重复。这正是 §4 要展开的「端到端为什么难」。
一句话记住:Flink 的 exactly-once = 内部状态一致(免费,靠 checkpoint)+ 对外副作用一致(收费,靠幂等或事务 sink)。
3. Flink 内部的 exactly-once:一致性快照 + 回滚 + 重放
先说「免费」的那一半。内部这一半的完整机制,状态管理与 Checkpoint 容错 会专门深挖,这里只把与 exactly-once 直接相关的三个环节快速串起来:一致性快照怎么切、失败怎么回滚、source 怎么重放。
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 条才是全部难点所在,§5 的两小节就分别讲幂等 sink 与事务 sink。
5. 两种 sink 路线:幂等 vs 事务
「让 sink 不产生重复副作用」听起来像一件事,实际是两套完全不同的思路:一条是认下重复但让它无害,另一条是干脆不让可能被回滚的数据露头。前者简单、延迟低,后者通用、代价是延迟。
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,恢复时按快照里的 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 之后,底层走的就是上一节那套 2PC:用 Kafka 的事务性 producer,每个 checkpoint 对应一个事务,preCommit 时 flush、commit 时提交。transactionalIdPrefix 是必填项——Flink 会基于它给每个并行 sink 子任务生成唯一的 transactional.id,Kafka 正是靠这个 id 做事务的 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 间隔「量化」了:间隔设 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 + 作业最坏恢复时间,并留足余量。
同时还有两个约束要一起看:broker 端的 transaction.max.timeout.ms(默认 15 分钟)是硬上限,producer 端设的值不能超过它,想要更大就得先调 broker;而 Flink 侧 Kafka producer 的默认事务超时是 1 小时,反而比这个上限还大,照抄默认值会被 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 一致性快照算法的理论源头。