Ch.11-12: 有界输入跑批、无界输入持续跑批 — 本质都是"确定性的纯函数 + 可重放的输入"; 开窗必须用事件时间
批处理像月末对账: 输入是一叠定死的账本 (有界、只读、不可变), 算错就改公式重算 — 大胆犯错是它的特权; 流处理是值班收银: 账本永远合不上 (无界), 但"确定性纯函数 + 可重放输入"的原则不变 — Kafka 的 offset 就是"这页账我看到哪了", 拨回去就能重算。最反直觉的是时间: 用户 10:00:59 发的帖子可能 10:03 才到达 — 按"到达时间"分桶它算 10:03 的量 (伪尖峰), 按"事件时间"分桶才是 10:00 的真相, watermark则是对天空喊话"10:02 前的都到齐了吧?" — 到齐才结账。
# 批: 处理 2026-09 的全部订单 → 输出月报 (会结束) # 流: 持续处理订单事件 → 实时大盘 (不结束)
# 经典词频: sort 承担了"shuffle+排序" cat logs | grep ERROR | awk '{print $5}' \ | sort | uniq -c | sort -rn | head
def mapper(line): # 无状态 for w in line.split(): emit(w, 1) def reducer(key, values): # 同 key 全到齐 emit(key, sum(values))
# reducer 3 拿到的: (dog,1)(dog,1)(dog,1) # 相邻 → 聚合一次线性扫完, 不用内存哈希
# reducer 看到 (按序): (user=42, PROFILE) # secondary sort 排前面 (user=42, ORDER#1) (user=42, ORDER#2) # → 一个 user 一次处理完, 不用缓存整个表
df = read.parquet(...) # 惰性: 只建计划 df2 = df.filter(...).groupBy(...).count() df2.write.parquet(...) # 触发执行 # 分区 3 丢了 → 按 lineage 只重算它, 不全量重跑
# 消费 = 移动指针, 不删数据 consumer.poll() # 从 offset 继续读 # 重放 = 把 offset 拨回去 (修投影的神器) consumer.seek(topic, offset=0)
# WAL 行: {"op":"u","table":"orders","after":{...}} # → Kafka topic "db.orders" → 各派生消费者 # 逻辑日志复制 (Ch.6) 与 CDC 是同一份流
T1: 写DB(A) ──────写MQ(A)
T2: 写DB(B) ──写MQ(B)
# MQ 顺序: A,B? B,A? 与 DB 提交序无关 → 分叉BEGIN; INSERT orders ...; INSERT outbox(topic, payload) VALUES('order-created', ...); COMMIT; # CDC 读 outbox → 发 Kafka → 删/标已发
# 昨天宕机 1h, 现在补 3600s 数据: # 按处理时间: "现在这秒" 3600 条 ← 伪尖峰 # 按事件时间: 均匀落回昨天各窗口 ← 真相
wm = max_event_ts - allowed_lateness(5min) # 窗口 [10:00,10:01) 在 wm≥10:01 时触发 # 10:00:59 的迟到事件 → retraction: 撤旧值发新值
# 恢复 = 从最近 checkpoint 重放输入 + 恢复状态 # sink 侧必须幂等 (UPSERT) 或事务写 — 否则重复
不用 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
核心思想: 用排序组织数据, 内存不够磁盘凑。
夜间作业把订单明细 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 作业, 中间不落盘。
下游发现 dt=0926 的宽表金额翻倍 — 上游重复文件, 按血缘只重算这一个分区。
# 排查: wide/dt=0926 的 lineage 上游 # orders/dt=0926 ← 上游采集重复写了一次 # 修复: 删掉重复源分区, 只重算受影响的下游 spark.read...filter(col("dt") == "2026-09-26").write.mode("overwrite")... # lineage 的价值: 故障半径 = 单分区, 不是整条管道
前提: 管道是确定性的 (同输入同输出), 不确定性函数要慎用。
搜索索引与库长期"偶发不一致", 抽查发现双写顺序在并发下相反。
# 事故态 (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。
投影代码有 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 页)。
一条畸形 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) — 混着处理两头堵。
按事件时间开 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 分钟, 窗口已发布 — 发撤销+修正, 报表自愈。
# 迟到事件落进已发布窗口 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 该调了。
库存服务重启要从头回放 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 管状态, 普通日志管事件流
这是"物化 = 重放日志"思想的存储层实现 (呼应事件溯源页)。
"框架说恰好一次"不够 — 端到端要生产侧事务 + 消费侧幂等同时成立。
# 生产侧: 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, 靠原子提交+幂等拼出来。
# 错: 改个统计口径 → 全表 8 小时 # 对: 按分区重算受影响的 dt=0926
# 错: reducer 里 dict 累积全部 key # 对: 依赖有序流线性聚合
# 错: groupByKey(celebrity_id) # 对: key+rand(0,9) 先聚, 再按原 key 二次聚
# 错: tmp 数据在 executor 本地盘 # 对: 中间结果 parquet 落 S3/HDFS
# 错: window(now, 60s) # 对: window(event_ts, 60s) + watermark
# 错: allowed_lateness = 0 / ∞ # 对: 5min + 迟到率监控
# 错: else: pass # 对: metrics.incr("late_dropped") + 告警
# 错: 风控 50ms 需求选 micro-batch # 对: batch=1s 的系统选 Flink/单事件流
# 错: checkpoint interval=100ms # 对: 30s + RocksDB 增量快照
# 错: "Flink exactly-once 所以不用管" # 对: DB UPSERT / 事务 sink
# 错: 处理一条扣一条 # 对: INSERT dedup(event_id) 再处理
# 错: save(db); kafka.send() # 对: 同事务写 outbox, CDC 发布
# 错: outbox 只插不删 # 对: published_at 分区, 7 天 TTL
# 错: 下游按位置解析 after 字段 # 对: 按名解析 + registry 兼容门禁
# 错: except: retry forever # 对: 3 次 → DLQ + 告警 + 重投工具
# 错: cleanup.policy=delete 永不设 TTL # 对: compact(状态) + delete retention(流)
# 错: enable.auto.commit=true # 对: 手动提交 + dedup(event_id)
# 错: order 流 × pay 流 全量缓存配对 # 对: join within 30min 窗口
# 错: "10:00 的算哪边?" 无答案 # 对: [10:00,10:01) + 边界事件测试用例
# 错: Java 批任务 + Python 流任务各写一遍 # 对: 同一 SQL/函数, 输入源不同