DDIA · 批处理与流处理

Ch.11-12: 有界输入跑批、无界输入持续跑批 — 本质都是"确定性的纯函数 + 可重放的输入"; 开窗必须用事件时间

MapReduce: shuffle 不是随洗牌, 是分布式排序 输入块 HDFS mapper 1 mapper 2 shuffle 按 key 重分布 hash(key) → reducer + 排序 (同 key 聚齐) reducer 1 reducer 2 mapper 无状态可并行; reducer 按 key 迭代 — 同 key 的所有记录 guaranteed 聚齐 任务级重试容错: 任何一个 task 失败重跑, 不影响其他 (输入只读不可变) sort-merge join: 两侧按 join key shuffle 排序 → reducer 一次处理一个 user 的全部记录 外部排序: 工作集超内存借磁盘顺序 IO (GNU sort 自动外排) — 用排序而不是内存哈希聚合 Spark: 整个 workflow 一个作业, 算子合并、中间态留内存; lineage 记依赖图, 分区丢按血缘重算 流处理: 事件时间 vs 处理时间 + 水位线 时间 窗口10:00 10:01 10:02 按 事件时间(event time) 分桶 straggler 迟到事件 事件到达 (乱序) watermark: 更早的不会再来了 → 触发 处理时间开窗的坑: 补积压时"到达时间"集中 → 伪尖峰; 必须按事件发生时间开窗 迟到策略: watermark 后丢弃+告警 / retraction 修正 (发一条撤销+新值) 窗口四类: tumbling 固定 / hopping 步进 / sliding 滑动 / session 活动聚团 micro-batching (Spark) ≈ 1s 攒批: 隐含 processing-time 窗口, 延迟下限 1s dual-write 之害 与 outbox + CDC 之解 dual-write (双写) 应用: 写 DB ✓ → 写 Kafka ✓ 并发交错: T1写库 T2写库 T2写MQ T1写MQ → 顺序反了, 索引与库永久不一致 outbox + CDC 业务写 + 事件写 同库同事务 → 事件表 (outbox) 落库 Debezium 从 WAL 读 outbox → 对外发布: 顺序 = 提交顺序 派生数据管道全景: 一条变更流喂饱所有视图 DB WAL → Debezium CDC → Kafka (log-based, offset 可重放) ├─→ 搜索索引 (Elasticsearch) ├─→ 缓存失效 ├─→ 物化视图 (流式 IVM) └─→ 数据湖 (T+1 分析) log compaction: 按 key 只留最新 → 新消费者从 0 扫描即得全量状态

批与流是同一件事

  • • 批: 有界输入 → 派生输出, 指标是吞吐
  • • 流: 无界输入持续跑批
  • • 共同点: 确定性纯函数 + 可重放输入
  • • 都可重跑 — 出错改代码重放即可

时间的纪律

  • • 开窗必须用事件时间, 不是处理时间
  • • watermark 宣告"到齐了"才触发
  • • 迟到: 丢弃+告警 或 retraction 修正
  • • processing-time 开窗 = 补积压伪尖峰

CDC 是派生的发动机

  • • dual-write 并发交错 → 永久不一致
  • • outbox 同事务 + CDC 从 WAL 发布
  • • offset 回拨 = 重放 = 修投影
  • • DLQ 收容毒消息, 别让它们卡管道

💡 一句话理解

批处理像月末对账: 输入是一叠定死的账本 (有界、只读、不可变), 算错就改公式重算 — 大胆犯错是它的特权; 流处理是值班收银: 账本永远合不上 (无界), 但"确定性纯函数 + 可重放输入"的原则不变 — Kafka 的 offset 就是"这页账我看到哪了", 拨回去就能重算。最反直觉的是时间: 用户 10:00:59 发的帖子可能 10:03 才到达 — 按"到达时间"分桶它算 10:03 的量 (伪尖峰), 按"事件时间"分桶才是 10:00 的真相, watermark则是对天空喊话"10:02 前的都到齐了吧?" — 到齐才结账。

🧠 必知必会 必考 & 必会

bounded / unbounded
批处理前提是有界数据 (文件/表), 流处理面对无界流 (持续事件)。同一套"纯函数+重放"思想, 输入形态不同。
# 批: 处理 2026-09 的全部订单 → 输出月报 (会结束)
# 流: 持续处理订单事件 → 实时大盘 (不结束)
Unix pipeline
小工具 stdin/stdout 组合; 用排序而非内存哈希聚合 — 工作集超内存也能算 (外排序), 是 MapReduce 的思想原型。
# 经典词频: sort 承担了"shuffle+排序"
cat logs | grep ERROR | awk '{print $5}' \
  | sort | uniq -c | sort -rn | head
MapReduce 流程
mapper (无状态并行) → shuffle (按 key 重分布+排序) → reducer (按 key 迭代); 任务级重试容错 — 输入只读不可变, 挂了重跑即可。
def mapper(line):    # 无状态
    for w in line.split(): emit(w, 1)
def reducer(key, values):   # 同 key 全到齐
    emit(key, sum(values))
shuffle 真相
shuffle 是分布式排序: mapper 按 key 哈希写给各 reducer + reducer 内排序 — "同 key 必然聚在同一个 reducer 且相邻"。
# reducer 3 拿到的: (dog,1)(dog,1)(dog,1)
# 相邻 → 聚合一次线性扫完, 不用内存哈希
sort-merge join
两侧按 join key shuffle 排序后, reducer 一次处理"一个 user 的全部记录"; secondary sort 保证"user 记录先于事件记录"到达。
# reducer 看到 (按序):
(user=42, PROFILE)    # secondary sort 排前面
(user=42, ORDER#1) (user=42, ORDER#2)
# → 一个 user 一次处理完, 不用缓存整个表
DAG / lineage / lazy
Spark 把整个工作流当一个作业: 算子合并、中间态留内存; lineage 记录算子依赖图, 分区丢失按血缘重算; lazy evaluation 先建计划优化后执行。
df = read.parquet(...)          # 惰性: 只建计划
df2 = df.filter(...).groupBy(...).count()
df2.write.parquet(...)          # 触发执行
# 分区 3 丢了 → 按 lineage 只重算它, 不全量重跑
log-based broker
Kafka: 持久追加日志 + 分区 + offset; 非破坏读 (读不删数据)、可重放、天然扇出 — 20TB 盘 @250MB/s ≈ 22 小时日志缓冲。
# 消费 = 移动指针, 不删数据
consumer.poll()          # 从 offset 继续读
# 重放 = 把 offset 拨回去 (修投影的神器)
consumer.seek(topic, offset=0)
CDC 变更数据捕获
从数据库 WAL 抽变更流 (Debezium/Kafka Connect) 喂给派生系统: 索引、缓存、物化视图、数仓 — 顺序 = 提交顺序, 天然一致。
# WAL 行: {"op":"u","table":"orders","after":{...}}
# → Kafka topic "db.orders" → 各派生消费者
# 逻辑日志复制 (Ch.6) 与 CDC 是同一份流
dual-write 之害
应用分别写库与索引: 并发交错导致永久不一致 (两个写者的顺序在两边相反) — 不是"概率小", 是结构必然。
T1: 写DB(A) ──────写MQ(A)
T2:     写DB(B) ──写MQ(B)
# MQ 顺序: A,B?  B,A?  与 DB 提交序无关 → 分叉
outbox pattern
业务写与事件写同库同事务 (先落 outbox 表), 由 CDC 从 WAL 读 outbox 对外发布 — 顺序与提交一致, 永不分叉。
BEGIN;
  INSERT orders ...;
  INSERT outbox(topic, payload) VALUES('order-created', ...);
COMMIT;   # CDC 读 outbox → 发 Kafka → 删/标已发
event / processing time
事件时间=事情发生时刻 (埋点携带), 处理时间=被处理时刻。开窗必须用前者 — 补积压时处理时间集中产生伪尖峰。
# 昨天宕机 1h, 现在补 3600s 数据:
# 按处理时间: "现在这秒" 3600 条 ← 伪尖峰
# 按事件时间: 均匀落回昨天各窗口 ← 真相
watermark 与 straggler
watermark 宣告"事件时间 < t 的不会再有", 据此触发窗口; 之后才到的 straggler: 丢弃+告警 或 retraction 修正 — 权衡要显式。
wm = max_event_ts - allowed_lateness(5min)
# 窗口 [10:00,10:01) 在 wm≥10:01 时触发
# 10:00:59 的迟到事件 → retraction: 撤旧值发新值
checkpoint 与 exactly-once
barrier 随流注入, 算子对齐后快照状态 (Flink); 端到端 exactly-once = 框架快照 + 原子提交/幂等 sink — 只快照不幂等仍会重复。
# 恢复 = 从最近 checkpoint 重放输入 + 恢复状态
# sink 侧必须幂等 (UPSERT) 或事务写 — 否则重复

🏭 生产实战 real world

场景 1 · 10GB 日志词频: Unix 管道 30 秒跑完

不用 Hadoop, 单机管道 + 外部排序就是最小的 MapReduce。

# 10GB nginx 日志统计 TOP URL (sort 自动外排)
cat access.log \
 | awk '{print $7}' \
 | grep -v '\.\(png\|css\|js\)' \
 | sort | uniq -c \
 | sort -rn | head -20

# map: awk 提取 → shuffle: sort → reduce: uniq -c
# 30s 跑完;  上规模 = 同结构换 Spark

核心思想: 用排序组织数据, 内存不够磁盘凑。

场景 2 · Spark DataFrame: 用户-订单宽表日更

夜间作业把订单明细 join 用户维表, 输出分析宽表到列存。

orders = spark.read.parquet("s3://lake/orders/dt=2026-09-26")
users  = spark.read.parquet("s3://lake/users")

wide = (orders.join(users, "user_id", "left")
        .groupBy("city", "category")
        .agg(sum("amount").alias("gmv"),
             countDistinct("user_id").alias("uv")))

wide.write.partitionBy("city").parquet("s3://lake/wide/")
# lineage: 失败重跑;  分区裁剪: 只读当天分区

对比 Hadoop: 同逻辑 3 个 MR 作业 → 1 个 Spark 作业, 中间不落盘。

场景 3 · lineage 排障: 输出分区错了只重算它

下游发现 dt=0926 的宽表金额翻倍 — 上游重复文件, 按血缘只重算这一个分区。

# 排查: wide/dt=0926 的 lineage 上游
#   orders/dt=0926 ← 上游采集重复写了一次

# 修复: 删掉重复源分区, 只重算受影响的下游
spark.read...filter(col("dt") == "2026-09-26").write.mode("overwrite")...
# lineage 的价值: 故障半径 = 单分区, 不是整条管道

前提: 管道是确定性的 (同输入同输出), 不确定性函数要慎用。

场景 4 · 双写不一致事故 → outbox + CDC 改造

搜索索引与库长期"偶发不一致", 抽查发现双写顺序在并发下相反。

# 事故态 (dual-write):
db.save(order)                    # 提交序: A 后 B
kafka.send(order)                 # 但 MQ 里 B 先 A 后
# → 索引最终状态 = A, 库最终状态 = B → 永久分叉

# 改造 (outbox + CDC):
BEGIN;
  upsert(orders, order);
  INSERT INTO outbox(agg_id, topic, payload) VALUES(order.id, ...);
COMMIT;
# Debezium 读 WAL 里的 outbox → Kafka (顺序=提交序)
# 索引消费者按序应用 → 与库最终一致 (永远)

上线后"索引不一致工单"从每周 3 单变 0。

场景 5 · offset 重放: 修复写错的物化视图

投影代码有 bug 算错了三天的推荐分数 — 把 offset 拨回 72 小时前重放。

# 1) 修 bug, 部署新版本消费者
# 2) 重置消费位点 (按时间戳)
kafka-consumer-groups --reset-offsets \
  --group rec-projection --topic events \
  --to-datetime 2026-09-23T00:00:00 --execute

# 3) 从该位置重放, 物化视图被重新计算 (幂等 sink)
# 4) 校验抽样一致后恢复对外读

这就是"读模型是派生数据"的流式版本 (呼应 CQRS 页)。

场景 6 · 毒消息卡管道: DLQ 兜底

一条畸形 payload 让消费者反复崩溃, 整个分区停摆 — 重试 N 次进 DLQ。

for attempt in range(3):
    try:
        process(msg)
        break
    except DeserializationError:      # 重试也没用的毒消息
        dlq.send(msg, reason=traceback)
        break
    except TransientError:
        sleep(0.1 * 2**attempt)       # 瞬时错误退避重试

# DLQ 巡检: 积压 > 0 告警, 人工修数后重投

区分: 瞬时错误 (重试) vs 结构错误 (DLQ) — 混着处理两头堵。

场景 7 · 实时成交大盘: 事件时间窗口 + watermark

按事件时间开 1 分钟窗统计 GMV, 允许 5 分钟迟到 — Flink 式逻辑 Python 演示。

state = {}                       # 窗口 → 累计金额
wm = 0                           # 水位线
LATENESS = 300                   # 容忍 5 分钟乱序

def on_event(ev):
    global wm
    wm = max(wm, ev.ts - LATENESS)          # 推进水位
    w = ev.ts // 60                          # 事件时间分桶
    state[w] = state.get(w, 0) + ev.amount
    for win in [k for k in state if (k+1)*60 <= wm]:
        emit_result(win, state.pop(win))     # 到齐才发布

大盘数字从此与离线对账一致 (都在事件时间轴上)。

场景 8 · 迟到修正: retraction 而不是硬丢

支付回调晚了 8 分钟, 窗口已发布 — 发撤销+修正, 报表自愈。

# 迟到事件落进已发布窗口 win=10:00:
emit_retraction(win="10:00", old=125_000)   # 撤销旧值
state["10:00"] += late_amount
emit_result("10:00", state["10:00"])         # 发布新值

# 下游 (物化视图/大盘) 按 key UPSERT → 自动收敛
# 全链路要求: sink 幂等 upsert, 不允许 append-only

取舍写进监控: retraction 率 > 0.5% 说明 allowed_lateness 该调了。

场景 9 · 日志压实: 新消费者 10 分钟重建全量状态

库存服务重启要从头回放 3 年事件? compacted topic 只留每 key 最新值。

# topic 配置: 按 key 压实 (每 key 永远留最新)
topic_config = {"cleanup.policy": "compact"}

# 事件: stock:S1=100 → 98 → 95 → 99 (只留 99)
# 新消费者 from offset=0: 30 分钟扫完 3 年
#   → 直接得到当前全量库存, 无需读库快照

# 分工: compacted 管状态, 普通日志管事件流

这是"物化 = 重放日志"思想的存储层实现 (呼应事件溯源页)。

场景 10 · 端到端 exactly-once: 事务写 + 幂等消费

"框架说恰好一次"不够 — 端到端要生产侧事务 + 消费侧幂等同时成立。

# 生产侧: Kafka 事务 (读-处理-写 在一个事务里)
producer.begin_transaction()
producer.send("output", result)
producer.send_transactional_offset()   # offset 与输出同事务
producer.commit_transaction()

# 消费侧: isolation.level=read_committed 只读已提交
# 下游 sink (DB): 幂等 UPSERT 兜底 — 三层缺一不可
INSERT INTO results(id, v) VALUES(...)
  ON CONFLICT (id) DO UPDATE SET v = EXCLUDED.v;

诚实口径: 端到端 exactly-once = effectively-once, 靠原子提交+幂等拼出来。

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 出错就全量重算上游 — 10TB 源数据重扫一遍救一个分区. 原因: 没有派生/血缘意识. 正解: 分区化 + lineage 局部重算。
# 错: 改个统计口径 → 全表 8 小时
# 对: 按分区重算受影响的 dt=0926
坑 2 · shuffle 当随机 — 不知道它是排序, 写出内存聚合. 原因: 名字误导. 正解: 相信"同 key 聚齐", 用 sort-merge 思维设计。
# 错: reducer 里 dict 累积全部 key
# 对: 依赖有序流线性聚合
坑 3 · 倾斜 key 打爆 reducer — 一个大 V 的百万事件挤一个任务. 原因: 没处理热点. 正解: salting 拆分 + 二次聚合。
# 错: groupByKey(celebrity_id)
# 对: key+rand(0,9) 先聚, 再按原 key 二次聚
坑 4 · 中间结果不落可靠存储 — 作业链 3 小时, 第 3 环失败从头再来. 原因: 中间态只在内存/本地盘. 正解: 关键中间态写 DFS/表格式。
# 错: tmp 数据在 executor 本地盘
# 对: 中间结果 parquet 落 S3/HDFS
坑 5 · 用处理时间开窗 — 补积压时大盘"伪尖峰". 原因: 时间语义没分清. 正解: 事件时间开窗 + watermark。
# 错: window(now, 60s)
# 对: window(event_ts, 60s) + watermark
坑 6 · watermark 设 0 或无限 — 设 0 迟到全丢; 设无穷窗口永不发布. 原因: 没有权衡意识. 正解: 按业务迟到分布定 (如 5min) + retraction 兜底。
# 错: allowed_lateness = 0 / ∞
# 对: 5min + 迟到率监控
坑 7 · 迟到静默丢弃 — 0.3% 的支付事件被扔, 对账永远差一点. 原因: 没有告警. 正解: 丢弃计数器告警 + 重要流 retraction。
# 错: else: pass
# 对: metrics.incr("late_dropped") + 告警
坑 8 · 微批当真流 — Spark Streaming 1s 攒批, 说好"毫秒级"延迟. 原因: 架构评审没问延迟下限. 正解: 毫秒级需求用真流 (Flink)。
# 错: 风控 50ms 需求选 micro-batch
# 对: batch=1s 的系统选 Flink/单事件流
坑 9 · checkpoint 太频繁 — 每 100ms 快照一次, 吞吐掉一半. 原因: 一刀切. 正解: 按恢复 RPO 定 (如 30s~1min), 大状态分开存。
# 错: checkpoint interval=100ms
# 对: 30s + RocksDB 增量快照
坑 10 · 框架内 exactly-once 当端到端 — Flink 内部恰好一次, sink 到 DB 照样重复. 原因: 只看了引擎宣传. 正解: sink 幂等/事务 + 输入可重放。
# 错: "Flink exactly-once 所以不用管"
# 对: DB UPSERT / 事务 sink
坑 11 · 消费无幂等 — rebalance 重投, 库里双倍库存扣减. 原因: at-least-once 语义没消化. 正解: 唯一键 UPSERT / 事件 id 去重表。
# 错: 处理一条扣一条
# 对: INSERT dedup(event_id) 再处理
坑 12 · dual-write 双写 — 库和索引并发交错永久分叉. 原因: 图省事两条都写. 正解: outbox + CDC 单一事实流。
# 错: save(db); kafka.send()
# 对: 同事务写 outbox, CDC 发布
坑 13 · outbox 不清理 — 发件箱表涨到 20 亿行, CDC 扫描越来越慢. 原因: 只管发不管收. 正解: 发布确认后清理/标记 + 分区裁剪。
# 错: outbox 只插不删
# 对: published_at 分区, 7 天 TTL
坑 14 · CDC 无视 schema 演进 — 上游加字段, Debezium 事件下游解析崩. 原因: 与编码页脱节. 正解: registry 校验 + 下游 unknown-tolerant。
# 错: 下游按位置解析 after 字段
# 对: 按名解析 + registry 兼容门禁
坑 15 · 无 DLQ — 一条毒消息卡死整个分区消费. 原因: 重试无分类. 正解: 结构性错误 N 次后进 DLQ + 巡检。
# 错: except: retry forever
# 对: 3 次 → DLQ + 告警 + 重投工具
坑 16 · 日志无限增长 — 事件 topic 三年 200TB, 重放=灾难. 原因: 没分"事件流"与"状态流". 正解: 状态用 compacted topic, 历史按 TTL/归档。
# 错: cleanup.policy=delete 永不设 TTL
# 对: compact(状态) + delete retention(流)
坑 17 · offset 语义混乱 — 先提交后处理, 崩了丢消息. 原因: 自动提交默认开. 正解: 处理完手动提交 + 消费幂等 (允许重复不许丢失)。
# 错: enable.auto.commit=true
# 对: 手动提交 + dedup(event_id)
坑 18 · join 无状态爆炸 — 流流 join 缓存所有左流事件等配对. 原因: 没设窗口边界. 正解: interval join 限定时间窗 + TTL 清理。
# 错: order 流 × pay 流 全量缓存配对
# 对: join within 30min 窗口
坑 19 · 窗口边界歧义 — 各实现左闭右开不一致, 对账差 1 条. 原因: 边界没写进约定. 正解: 文档声明 [start, end), 边界事件单测覆盖。
# 错: "10:00 的算哪边?" 无答案
# 对: [10:00,10:01) + 边界事件测试用例
坑 20 · 批流两套代码 — lambda 式双实现, 口径永远对不齐 (Ch.13 展开). 原因: 历史包袱. 正解: 同一套确定性逻辑跑在批/流两输入上 (Beam/kappa)。
# 错: Java 批任务 + Python 流任务各写一遍
# 对: 同一 SQL/函数, 输入源不同