Skip to content
Charles Shao
Go back

Flink 深挖 · 端到端 Exactly-Once 与两阶段提交

views

在广告实时计费里,一条曝光被处理两次,就是账上多算了一笔钱——多计费轻则对账对不平、财务月底追着你要说法,重则客户投诉超扣、平台被迫赔付。所以”作业重启后数据别重复”这件事,在计费/结算链路上不是”锦上添花”,而是红线。而这条红线的技术名字,就叫 exactly-once(精确一次)

问题是,“精确一次”这四个字被用得太随意,很多人以为它是”每条消息只被处理一次”——这是个天大的误解,也是本篇要拆开的第一件事。

一句话定位:exactly-once 不是”每条消息只被物理处理一次”,而是”故障重放后,反映到状态与对外结果上等效于一次”(effectively-once);端到端做到它 = 可重放的 source + 一致性快照 + 事务/幂等的 sink,三者缺一不可。

TL;DR

Table of contents

Open Table of contents

1. 先把三种投递语义说清楚

任何”可靠性”讨论都得先定义”可靠到什么程度”。流处理里有三个层级,区别只在故障重放时会发生什么

三种投递语义在故障重放后的结果对比图。三列并排:左列 at-most-once(至多一次,≤1)——机制是发出即不管、不重放不重试,故障时处理到一半的数据直接丢失,结果是可能丢数据,落到广告计费就是漏计费导致平台少收钱;中列 at-least-once(至少一次,≥1)——机制是 source 可重放但写出可能重复,故障时已处理的数据被重放、重复生效,结果是可能重复,落到计费就是重复计费导致对账不平;右列 exactly-once(精确一次/状态一致)——机制是一致性快照加回滚加事务提交,故障时回滚到快照后重放但对外只生效一次,结果是不多不少正好一次、计费精确可对账。底部强调:exactly-once 不是每条消息只被物理处理一次,而是反映到状态与对外结果上等效于一次(effectively-once),失败时消息仍会被重放、算子仍会重复计算,但靠一致性快照回滚加事务或幂等写,最终对外只生效一次。

三种语义的分水岭都在”故障重放时会怎样”:at-most-once 丢、at-least-once 重复、exactly-once 靠快照回滚 + 事务/幂等做到”效果上一次”。

落到广告实时计费上,为什么必须是最贵的那一档?因为另外两档的代价都是真金白银

计费是”重复即金钱损失、丢失即收入损失”的典型场景,两个方向都不能容忍——这就是为什么这条链路值得付出 exactly-once 的全部复杂度。

2. 澄清”精确一次”:它说的是状态,不是物理次数

这是全篇最重要的一个观念,先钉死它,后面才不会绕晕。

“exactly-once”并不意味着每条消息只被物理地处理一次。 恰恰相反:一旦发生故障,Flink 会把作业回滚到上一个一致性快照,然后从那个位点重放 source——这意味着快照之后、故障之前处理过的那批消息,会被再处理一遍。算子的 mapprocess、聚合逻辑,统统会对这批消息重复执行

那”精确一次”到底精确在哪?精确在最终状态与对外产生的结果上。业界更准确的叫法是 effectively-once(效果上一次)

不管中间重放了多少次、算子重算了多少遍,反映到 Flink 的托管状态和外部 sink 上的净效果,等同于每条消息恰好被处理了一次。

理解这一点,才能理解后面所有机制的设计动机:

一句话记住:Flink 的 exactly-once = 内部状态一致(免费,靠 checkpoint)+ 对外副作用一致(收费,靠幂等或事务 sink)。

内部这一半的完整机制,状态管理与 Checkpoint 容错 会专门深挖,这里只快速串起与 exactly-once 直接相关的三个环节。

3.1 一致性快照:异步屏障快照(ABS)

Flink 的 checkpoint 基于 Chandy-Lamport 分布式快照算法的变体——异步屏障快照(Asynchronous Barrier Snapshotting)。核心是往数据流里周期性注入一种特殊记录 checkpoint barrier

所有算子都完成快照 n 后,checkpoint n 才全局完成。这份全局快照的语义是:“如果从这里恢复,就像所有算子都恰好处理到 barrier n 之前的最后一条数据”。

3.2 失败回滚 + source 重放:一致性从何而来

当某个 task 失败,Flink 的恢复动作是整体性的(默认 region/full failover):

  1. 回滚状态:把所有算子的状态恢复到最近一次成功的 checkpoint n
  2. 回退 source 位点:checkpoint n 里存了 source 的读取位点(如 Kafka 的 offset)。恢复时 source 从这个 offset 重新开始读
  3. 重放: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 能回滚自己的状态,却回滚不了已经产生的外部副作用。 内部的一致性快照只解决了”我自己的账对不对”,解决不了”我已经泼出去的水”。

要补上这一半,只有两条路:

  1. 让 source 可重放(第 3 节已解决)——保证重放的数据源头是确定、可回退的。
  2. 让 sink 不产生重复副作用——要么写入本身幂等(重复写 = 写一次),要么把”写”延迟到”确定不会回滚”之后(事务)。

这两条路里第 2 条是全部难点所在,下面两节分别讲幂等 sink 和事务 sink。

5. 两种 sink 路线:幂等 vs 事务

5.1 幂等 sink:让”重复写”等于”写一次”

如果外部写入天然幂等——同样的数据写 N 次,效果和写 1 次一样——那重放导致的重复写就无害了。典型手段:

// 幂等写:按 (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 的通用骨架抽成了 TwoPhaseCommitSinkFunctionKafkaSinkEXACTLY_ONCE 模式内部就是这套逻辑)。它要求实现四个方法,每个都严格绑定 checkpoint 生命周期的一个时点

方法触发时机干什么
beginTransaction()上一个事务提交后 / 作业启动开启一个新事务,后续处理结果都写进它
preCommit()checkpoint 快照时snapshotStateflush 数据、把事务句柄写进 Flink 状态;进入”预提交”,但还不 commit
commit()checkpoint 全局完成通知时notifyCheckpointComplete真正 commit 事务,数据对下游可见
abort()失败 / 恢复时回滚事务,丢弃未提交的写入

把它和 checkpoint 生命周期对齐,时序是这样的:

两阶段提交与 Checkpoint 生命周期绑定的时序图。三个参与方从左到右:Checkpoint 协调者(JobManager)、事务 Sink 算子(TwoPhaseCommitSink)、外部系统(Kafka 事务)。步骤依次为:1) Sink 算子 beginTransaction 开启事务 Txn-n;2) Sink 把处理结果写入事务(未提交、下游看不到)到外部系统;3) 协调者向 Sink 注入 Checkpoint Barrier n;4) Sink 执行 preCommit:flush 数据、把 offset 与事务句柄写进快照;5) Sink 向协调者回复快照完成 ack(含待提交事务元数据,虚线);随后蓝色说明:所有算子都 ack 后 Checkpoint n 全局完成、持久化到状态后端;6) 协调者向 Sink 发送 notifyCheckpointComplete(n);7) Sink 向外部系统执行 commit:提交事务 Txn-n,数据对下游可见;8) Sink beginTransaction 开启下一个事务 Txn-(n+1)。底部红色说明:若任一步失败,abort 当前事务并回滚到上一个成功的 checkpoint、重放 source,外部只有已 commit 的数据可见(下游需 read_committed),未提交的写入被丢弃,故不重复。

两阶段提交的”提交点”与 checkpoint 完全对齐:checkpoint 快照时 preCommit(第一阶段),checkpoint 全局完成通知时才 commit(第二阶段)。这保证”只有不会被回滚的数据才对外可见”。

理解这套协议的关键,是看清 commit 为什么必须发生在 notifyCheckpointComplete

还有个容错细节:恢复时对”已 preCommit 但还没收到 commit 通知”的事务,要重新 commit(而不是 abort)。因为它对应的 checkpoint 已经成功了(否则不会恢复到它之后),这批数据是该生效的。TwoPhaseCommitSinkFunction 把待提交事务的句柄存进了 Flink 状态,恢复时能拿回来重新 commit——这就是为什么 preCommit 要”把事务句柄写进快照”。

7. Kafka 端到端 exactly-once:实战

把上面的理论落到最常见的 Kafka → Flink → Kafka 链路。这也是 Kafka 可靠性与 exactly-once 在 Flink 侧的落地。

Kafka 到 Flink 到 Kafka 的端到端 exactly-once 链路图。从左到右四个环节横向串联:KafkaSource(offset 存进 checkpoint)—可重放—> Flink 算子(状态随 checkpoint 快照)—两阶段提交—> KafkaSink(EXACTLY_ONCE / txn.id)—只读已提交—> 下游消费者(read_committed)。下方黄色说明:端到端一致等于三件事缺一不可:① source 可重放,offset 随 checkpoint 一起快照、回滚时一并回退;② sink 幂等或事务,KafkaSink 用事务性 producer、commit 与 checkpoint 完成对齐;③ 下游 isolation.level=read_committed,只读已提交的事务数据,未提交的写入对它不可见。再下方红色警告:关键坑——transaction.timeout.ms 必须大于 checkpoint interval 加作业恢复时间,否则事务在 commit 前被 broker 判超时中止,那段数据丢失。

端到端三要素:可重放 source(offset 进 checkpoint)+ 事务 sink(两阶段提交对齐 checkpoint)+ 下游 read_committed。任何一环缺失,端到端 exactly-once 就破。

7.1 KafkaSource:offset 存进 checkpoint

新版 KafkaSourceflink-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

KafkaSinkDeliveryGuarantee.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 与低延迟的根本矛盾

8.2 transaction.timeout.ms:最容易踩的雷

Kafka 的事务有超时时间 transaction.timeout.ms一个事务开启后,超过这个时间还没 commit,broker 会主动把它中止(abort)。

在 Flink 事务 sink 里,一个事务的存活时间是”从 beginTransactioncommit”,而 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,下游按它去重

这套幂等去重的更一般讨论(幂等键设计、去重窗口、最终一致),见 幂等与一致性

参考


views
Share this post on:

Previous Post
Flink 深挖 · DataStream API:ProcessFunction、状态、Timer 与 Side Output
Next Post
Flink 深挖 · 时间语义、窗口与 Watermark