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,于是快照之后、故障之前处理过的那批消息会被再处理一遍,算子里的 map、process 与聚合逻辑统统对它们重复执行一轮。

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

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

把这句话理解透,后面所有机制的设计动机就都顺了:

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

先说「免费」的那一半。内部这一半的完整机制,状态管理与 Checkpoint 容错 会专门深挖,这里只把与 exactly-once 直接相关的三个环节快速串起来:一致性快照怎么切、失败怎么回滚、source 怎么重放。

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 条才是全部难点所在,§5 的两小节就分别讲幂等 sink 与事务 sink。

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

「让 sink 不产生重复副作用」听起来像一件事,实际是两套完全不同的思路:一条是认下重复但让它无害,另一条是干脆不让可能被回滚的数据露头。前者简单、延迟低,后者通用、代价是延迟。

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

方法触发时机干什么
beginTransaction()上一个事务提交后 / 作业启动开启一个新事务,后续处理结果都写进它
preCommit()checkpoint 快照时(snapshotState)flush 数据、把事务句柄写进 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

新版 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 与低延迟之间的根本矛盾,只能按业务容忍度取舍:

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,让下游按它去重。

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

参考


–views
Share this post on:

Previous Post
Flink 深挖 · DataStream API 与有状态处理
Next Post
Flink 深挖 · 时间语义、窗口与 Watermark