Streaming 精确一次语义
概述
实时计算的可靠性核心是交付语义:At-Least-Once、At-Most-Once 还是 Exactly-Once。本文讲清三种语义的定义,checkpoint/WAL 如何保证计算不丢,幂等与事务输出如何保证结果不重,最终拼出端到端精确一次。
一、三种交付语义
| 语义 | 含义 | 风险 |
|---|---|---|
| At-Most-Once | 至多一次,可能丢 | 数据丢失 |
| At-Least-Once | 至少一次,可能重 | 数据重复 |
| Exactly-Once | 恰好一次 | 实现复杂 |
1.1 为什么难
流处理要精确一次,必须三端配合:
源(Source) 计算(引擎) 输出(Sink)
可重放 可恢复状态 幂等或事务缺任何一环都会出现丢或重。
二、计算侧:checkpoint 与 WAL
2.1 Checkpoint
| 内容 | 说明 |
|---|---|
| 状态 | 聚合中间状态 |
| Offset | 消费位置 |
| 元数据 | 配置与算子定义 |
每批执行前/后写 checkpoint
故障后从 checkpoint 恢复
→ 计算结果不丢(可重放)| 参数 | 说明 |
|---|---|
spark.streaming.checkpoint.interval | 批量 checkpoint 间隔 |
| checkpointLocation | HDFS 可靠存储路径 |
2.2 WAL(Write-Ahead Log)
Spark Streaming 的 WAL 把 Receiver 接收的数据先写日志再消费:
数据到达 → 写 WAL(HDFS) → 处理
Receiver 崩溃 → 从 WAL 重放| 对比 | 无 WAL | 有 WAL |
|---|---|---|
| 数据落盘 | 仅内存 | 日志先落 |
| 崩溃恢复 | 可能丢 | 可重放 |
| 性能 | 快 | 慢一些 |
scala
// 开启 WAL
sparkConf.set("spark.streaming.receiver.writeAheadLog.enable", "true")三、输出侧:幂等输出
3.1 原理
重复执行的结果可重复写入而结果一致。
3.2 常见幂等实现
| 方式 | 说明 |
|---|---|
| 文件覆盖 | 写固定分区目录,重复写覆盖 |
| 唯一键 | MySQL 唯一索引,冲突忽略 |
| 去重表 | 落库前查重(代价高) |
| Redis SETNX | 键存在即跳过 |
sql
-- MySQL 幂等写:唯一键 + INSERT IGNORE
INSERT IGNORE INTO dws_order_daily (dt, order_id, amount)
VALUES (?, ?, ?)| 手段 | 适用 |
|---|---|
| 文件 Sink | HDFS/Parquet(天然幂等) |
| 唯一键 | MySQL/PG |
| 状态去重 | 引擎内部 |
四、输出侧:事务输出
4.1 原理
把一批结果放在同一事务提交,要么全部生效,要么全部回滚:
事务写流程:
1. 批处理完成
2. 开启事务
3. 写入本批全部结果
4. 提交(或回滚)
5. 标记该批 offset 已提交4.2 两阶段提交思路
| 阶段 | 动作 |
|---|---|
| 预提交 | 结果写暂存,记录批次 ID |
| 确认 | offset 提交后,正式提交结果 |
Kafka 事务 + 引擎协调可实现端到端精确一次,但复杂度高,多数场景用幂等替代。
五、端到端精确一次方案
5.1 方案组合
| 端 | 策略 |
|---|---|
| 源 | Kafka 可重放(按 offset) |
| 计算 | checkpoint 恢复状态 |
| 输出 | 幂等(唯一键/文件覆盖)或事务 |
5.2 具体做法(Structured Streaming + Kafka)
scala
val query = df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("topic", "result")
.option("checkpointLocation", "hdfs:///ckpt/result")
.outputMode("append")
.start()| 组件 | 保证 |
|---|---|
| Kafka 源 | 按 checkpoint 记录 offset 重放 |
| 引擎 | 状态从 checkpoint 恢复 |
| Kafka Sink | 写入幂等(同 offset 同内容) |
| 结果 | 端到端精确一次 |
5.3 简化策略:幂等兜底
多数生产场景采用 At-Least-Once + 幂等输出,效果等价于精确一次且实现简单:
至少一次(可能重复)
+ 输出幂等(重复写入结果一致)
= 等效精确一次六、语义选型
| 场景 | 推荐 |
|---|---|
| 报表/聚合(可重跑) | At-Least-Once + 幂等 |
| 金额/订单(不可重复) | 端到端精确一次或幂等+唯一键 |
| 监控指标 | At-Least-Once 即可 |
| 低延迟丢弃容忍 | At-Most-Once(罕见) |
| 决策点 | 说明 |
|---|---|
| 业务容忍重复? | 容忍→At-Least-Once |
| 输出可幂等? | 可→幂等方案最简单 |
| 需要严格一次? | 事务 + 协调(复杂) |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 结果重复 | 输出无幂等,加唯一键或文件覆盖 |
| 恢复后数据重复 | 语义为至少一次,靠输出幂等兜底 |
| checkpoint 过大 | 状态大,设置超时清理或调 checkpoint 间隔 |
| WAL 太慢 | WAL 写 HDFS 有开销,权衡可靠性与性能 |
| 事务失败一半 | 需两阶段提交或改用幂等方案 |