引言
“我们用了 Flink,所以是精确一次”——这是流处理里最常见的误解。精确一次(Exactly-Once)从来不是某个引擎的开关,而是源、处理、汇三方协作才能达成的端到端保证。只把处理层配置对,源会重复发、汇会重复写,结果照样错。
Exactly-Once 是一条跨越整个链路的一致性协议,引擎只负责其中一环。
本文拆解这条链路的每一环:Checkpoint 如何对齐、状态如何持久化、两阶段提交如何保证汇端不重复、故障后如何恢复,以及生产中真正会踩的坑。
一、精确一次到底意味着什么
1.1 三种投递语义
| 语义 | 含义 | 典型场景 |
|---|---|---|
| At-Most-Once | 最多一次,可能丢 | 非关键日志 |
| At-Least-Once | 至少一次,可能重 | 可幂等去重的场景 |
| Exactly-Once | 恰好一次,不丢不重 | 计费、对账 |
Exactly-Once 的严格定义是:在故障恢复后,每条输入对状态与输出的影响恰好一次。注意它约束的是"影响",而非"物理投递次数"。
1.2 精确一次是端到端概念
# 端到端精确一次 = 源可回放 + 处理状态可恢复 + 汇可幂等/事务
# 缺任何一环, 都退化为 At-Least-Once
# 常见错觉: 只配了处理层, 源和汇没配合
- 源:必须支持按位点回放(Kafka offset、文件位置)。
- 处理:状态随 Checkpoint 一致地持久化。
- 汇:写入必须幂等或事务化。
1.3 为什么这么难
因为分布式系统里没有"全局时钟"。处理进程、状态存储、外部系统各自独立,要让它们在同一时刻达成一致,需要一套协调协议——这就是 Checkpoint 与两阶段提交存在的理由。
二、Flink 的 Checkpoint 机制
2.1 Barrier 与快照
Flink 周期性向数据流中注入 Barrier,Barrier 随数据流动,把流切分成"快照前"与"快照后"。算子收到所有输入的 Barrier 后,把自己的状态快照上报,形成一个全局一致的快照。
# Checkpoint 流程
# JobManager 触发 → 源注入 Barrier(n) → Barrier 随数据流动
# → 算子对齐 Barrier → 快照本地状态 → 上报完成
# → 全部完成 → Checkpoint n 完成(可用于恢复)
2.2 Barrier 对齐
对齐(Aligned):算子等所有输入的 Barrier 到齐再快照,保证一致性,但慢输入会拖慢整体。非对齐(Unaligned):不等待,把在途数据也存入快照,延迟更低,但快照更大。
# 对齐与非对齐 checkpoint 配置
execution.checkpointing:
interval: 60s
mode: exactly_once # 或 at_least_once
unaligned: false # true 则启用非对齐
timeout: 10min
min-pause: 30s
max-concurrent: 1
externalized:
enabled: true
retention: 3d # 保留外部快照, 支持手动恢复
2.3 Checkpoint 与 Savepoint
| 类型 | 触发者 | 用途 | 特点 |
|---|---|---|---|
| Checkpoint | 系统周期 | 故障自动恢复 | 轻量、可覆盖 |
| Savepoint | 用户手动 | 升级/迁移/调整 | 稳定、可移植 |
Savepoint 用于计划内操作(如改并行度、升级版本),Checkpoint 用于计划外故障恢复。
三、状态后端与状态存储
3.1 状态后端类型
| 后端 | 状态存放 | 适用 |
|---|---|---|
| HashMapStateBackend | JVM 堆内存 | 小状态、低延迟 |
| EmbeddedRocksDBStateBackend | 本地 RocksDB | 大状态、超内存 |
RocksDB 后端把状态放在本地磁盘的 LSM 结构里,支持超出内存的大状态,代价是读写有序列化开销。
# RocksDB 状态后端配置
state.backend: rocksdb
state.backend.incremental: true # 增量 checkpoint, 大幅降低开销
state.checkpoints.dir: s3://lake/flink/checkpoints
state.savepoints.dir: s3://lake/flink/savepoints
3.2 增量 Checkpoint
全量 Checkpoint 每次都把整个状态传到远端,大状态下代价极高。增量 Checkpoint 基于 RocksDB 的 SST 文件,只上传变化的部分,显著降低成本。
# 全量 vs 增量 checkpoint
# 全量: 状态 100GB → 每次上传 100GB, 分钟级
# 增量: 只上传新增 SST → 通常几 GB, 秒级
# 前提: 使用 RocksDB 后端 + 开启 incremental
3.3 状态结构设计
状态不是越多越好。用KeyedState(ValueState/ListState/MapState)而非算子状态,才能随 key 分区并行;避免无界增长的 ListState,必要时设 TTL。
// 带 TTL 的 ValueState, 防止状态无限增长
ValueStateDescriptor<Long> desc =
new ValueStateDescriptor<>("lastSeen", Long.class);
StateTtlConfig ttl = StateTtlConfig.newBuilder(Duration.ofHours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
desc.enableTimeToLive(ttl);
四、两阶段提交与端到端精确一次
4.1 两阶段提交 Sink
要让汇端"不重复写",Sink 必须参与 Checkpoint:数据先写入预提交状态,Checkpoint 完成时才真正提交,失败则回滚。这就是两阶段提交(2PC)。
# Flink 2PC Sink 流程
# 1. preCommit: 数据写入事务/临时区(如 Kafka 未提交事务)
# 2. Checkpoint 完成: 提交事务(如 commitOffsets)
# 3. Checkpoint 失败: 回滚, 丢弃未提交数据
4.2 Kafka Sink 的事务实现
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("orders-out")
.setValueSerializationSchema(new SimpleStringSchema())
.build())
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("flink-orders")
.setProperty("transaction.timeout.ms", "900000")
.build();
transactionalIdPrefix 用于生成唯一事务 ID,transaction.timeout.ms 必须大于 Checkpoint 间隔与最大暂停之和,否则事务会被 broker 提前中止。
4.3 幂等 Sink 作为替代
并非所有 Sink 都支持事务。对于不支持 2PC 的系统(如某些 OLAP、对象存储),可用幂等写替代:以确定性主键做 Upsert,重复写入结果相同。这要求下游支持主键去重。
五、故障恢复与状态一致性
5.1 恢复流程
# 故障恢复
# 1. 从最近完成的 Checkpoint 读取状态
# 2. 源按快照记录的位点回放(如 Kafka offset)
# 3. 重放位点之后的数据, 状态被重建
# 4. 未提交的 Sink 事务被回滚或忽略
5.2 恢复的一致性级别
| 配置 | 恢复后一致性 | 代价 |
|---|---|---|
| exactly_once + 对齐 | 精确一次 | 延迟受慢流影响 |
| exactly_once + 非对齐 | 精确一次 | 快照更大 |
| at_least_once | 至少一次 | 需下游去重 |
5.3 恢复时间与状态大小
恢复时间 ≈ 状态大小 / 恢复带宽。状态越大,恢复越慢,RTO 越长。这是"状态治理"重要的根本原因:不仅影响运行时,更影响故障时的恢复速度。
六、状态大小治理与调优
6.1 状态膨胀的来源
- 无界聚合:按 user_id 累加但 key 无限增长,状态永不清理。
- 超长窗口:窗口过大,缓存数据多。
- 未设 TTL:临时状态一直保留。
- 大对象状态:把整条记录塞进状态,而非只存必要字段。
6.2 治理手段
# [ ] 为状态设置 TTL, 定期清理过期 key
# [ ] 用增量聚合代替缓存全量(如 reduce/aggregate)
# [ ] 控制 key 基数, 避免高基数维度
# [ ] 状态只存必要字段, 不存整条记录
# [ ] 合理设置并行度, 均衡 key 分布
# [ ] 监控状态大小与 checkpoint 耗时趋势
6.3 RocksDB 调优
# RocksDB 关键调优
state.backend.rocksdb.memory.managed: true # 托管内存, 防 OOM
state.backend.rocksdb.block.cache-size: 256mb
state.backend.rocksdb.writebuffer.size: 64mb
state.backend.rocksdb.compaction.style: LEVELED
七、Kafka/Flink 精确一次实战
7.1 完整链路配置
# Kafka source → Flink 处理 → Kafka sink 精确一次
# 1. source: 记录 offset 到 checkpoint
# 2. 处理: 状态随 checkpoint 持久化
# 3. sink: 事务提交与 checkpoint 对齐
# 4. 消费者: isolation.level=read_committed
-- Flink SQL 端到端精确一次
SET 'execution.checkpointing.interval' = '60s';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
SET 'execution.checkpointing.timeout' = '10min';
CREATE TABLE orders_src (
order_id STRING, user_id BIGINT, amount DOUBLE,
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.isolation.level' = 'read_committed',
'scan.startup.mode' = 'group-offsets',
'format' = 'json'
);
CREATE TABLE orders_sink (
user_id BIGINT, total DOUBLE
) WITH (
'connector' = 'kafka',
'topic' = 'orders_agg',
'properties.bootstrap.servers' = 'kafka:9092',
'sink.delivery-guarantee' = 'exactly-once',
'sink.transactional-id-prefix' = 'flink-agg',
'format' = 'json'
);
INSERT INTO orders_sink
SELECT user_id, sum(amount) FROM orders_src GROUP BY user_id;
7.2 验证精确一次
# 验证方法
# 1. 注入已知条数的测试数据
# 2. 人为 kill TaskManager 触发恢复
# 3. 校验汇端结果条数 == 输入条数, 且无重复
# 4. 检查下游是否读到未提交数据(read_committed)
7.3 常见失败模式
- 事务超时:
transaction.timeout.ms小于 checkpoint 间隔,事务被中止,任务失败。 - 消费者读未提交:
isolation.level未设为read_committed,读到脏数据。 - 位点未提交:source 未启用 checkpoint 记录 offset,恢复后重复消费。
八、踩坑清单与最佳实践
8.1 配置类坑
- Checkpoint 间隔过大:恢复要回放更多数据,RTO 长。
- 超时过小:大状态快照来不及,任务反复失败。
- 未开外部化:任务取消后 Checkpoint 丢失,无法恢复。
- 并行度改后无 Savepoint:直接从 Checkpoint 恢复可能状态不兼容。
8.2 状态类坑
- 无界状态:key 基数爆炸,状态撑爆磁盘。
- 未设 TTL:临时状态永久驻留。
- RocksDB 未托管内存:与算子抢内存导致 OOM。
8.3 Sink 类坑
- 假精确一次:只配了处理层,Sink 是普通 At-Least-Once,端到端仍是重复。
- 幂等键不确定:重放时主键变化,幂等失效。
- 下游不支持事务:硬套 2PC 失败,应改用幂等写。
总结
| 环节 | 精确一次的保证手段 | 关键配置 |
|---|---|---|
| 源 | 位点随快照记录、可回放 | 记录 offset |
| 处理 | Barrier 对齐、状态快照 | exactly_once 模式 |
| 状态 | 增量快照、TTL 治理 | RocksDB + incremental |
| 汇 | 两阶段提交或幂等写 | delivery-guarantee |
| 恢复 | 从快照重建 + 位点回放 | externalized 快照 |
| 消费 | 只读已提交数据 | read_committed |
精确一次是一条跨越源、处理、汇的端到端协议,任何一环掉链子都会退化为至少一次。真正落地时,重点不在于"打开某个开关",而在于:状态要治理(否则恢复慢、磁盘爆)、Sink 要真正事务化或幂等(否则重复写)、消费端要只读已提交(否则读到脏数据)。理解这三件事,你才真正掌握了流处理的精确一次。
参考与延伸阅读
- Apache Flink 官方文档:Checkpointing、State Backends 与 Exactly-Once
- Apache Kafka 官方文档:事务与 read_committed 语义
- Flink 官方博客:End-to-End Exactly-Once 实现原理
- Apache Flink 流处理 — 流处理引擎基础
- 流处理窗口与时间语义 — 时间与窗口机制
- Kafka Connect 与 CDC — 源端变更捕获
- 实时数据仓库 — 精确一次在实时链路中的位置
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。