跳到主要内容

L2.9 流式 ETL 与 Embedding 管道

三维坐标 layer: L2(数据与算子层)level: Seniorpillar: 数据

L2.1 以批 ETL 与 DVC 为主;L2.2 讲向量库。本章补齐 流式路径Embedding 工业化管道——RAG 知识库如何从「 nightly 全量重建」升级为「增量近实时 + 可控 backfill」。

学习目标

  • 前置知识:读过 L2.1、L2.2;知道 Kafka 是消息队列即可。无需自建 Flink 集群(可用 Docker Compose 或概念+脚本演练)。
  • 学完产出:① 能画出 Kafka → Flink → Sink 的流式链路,说清 event-time 窗口与 watermark;② 能解释 Flink exactly-once 与 Kafka 事务的关系;③ 能设计 Embedding 三模式:batch 全量 / streaming 增量 / backfill 补历史;④ 能列出 embedding 版本变更时的「双写 + 切换 + 旧索引 TTL」步骤;⑤ 用 Python 模拟 mini streaming pipeline(queue + worker)跑通增量写入向量库 mock。
  • 阅读姿势:RAG 线上事故半数来自 索引与 embedding 模型版本不同步——本章把这类变更做成可运维流程。

背景与现状

路径延迟典型用途
批 ETL(Spark)小时~天预训练语料、离线评估集
流 ETL(Kafka+Flink)秒~分钟业务日志、实时特征、文档变更事件
Embedding batch小时初次建库、模型大版本升级
Embedding streaming分钟新知识入库、工单/文档更新

📅 时效说明(2026-07 口径):消息层主线已进入 Kafka 4.x 线——KRaft 共识全面替代 ZooKeeper(4.x 线已彻底移除 ZooKeeper 依赖),单进程即可承担 broker + controller 双角色;计算层主线为 Flink 2.x 线(存算分离的 disaggregated state、更成熟的流批一体)。本文示例 Compose 采用的正是 KRaft 单节点模式,与 4.x 线配置同构,替换镜像 tag 即可平滑演练。版本细节随发布演进,动手前以两个项目的官方 release notes 为准。

两个真实业务症状(带着问题读)

  • 症状 1:「新文档搜不到」工单激增,检索延迟却一切正常。某企业知识库 RAG 的 Grafana 大盘显示查询 P99 稳定,但用户抱怨上午发布的公告下午还搜不到。排查发现 embedding 消费组的 Kafka lag 已积压约 2 小时,且斜率为正(只涨不落)。根因:前一晚 embedding 模型从 base 升级到 large,单批推理耗时翻倍,消费吞吐从 640 docs/s 掉到 320 docs/s,跌破日间约 400 docs/s 的稳态生产速率——这是典型的背压(backpressure)失衡,思考题 2 会把这笔账完整算一遍。
  • 症状 2:同一段落在检索结果里出现 3 次。夜间 Flink job OOM 重启后,at-least-once 语义把上一个 checkpoint 之后的 doc.updated 事件全部重放;而向量库 Sink 写的是 insert 而非按确定性主键 upsert,重复向量把 top-10 的有效信息位挤掉,检索质量静默下降。思考题 1 讨论幂等设计如何让「重放」变得无害。

企业 RAG 常见架构:

原理与架构

概念含义踩坑
Event Time事件真实发生时间用 processing time 乱序严重
Watermark允许迟到数据的边界过小丢数据,过大延迟高
Windowtumbling / sliding / session无界聚合必须加窗口
CheckpointFlink 状态快照间隔影响恢复粒度
Exactly-once端到端不丢不重需 Kafka 事务 + 两阶段提交 Sink

投递语义的三档梯子(决定故障重启后世界长什么样):

  • at-most-once:收到即提交 offset,处理失败就丢——RAG 场景等于「文档静默缺失」,基本不用。
  • at-least-once:先处理后提交,失败重放——不丢但会重。Flink 从 checkpoint 恢复时,checkpoint 之后已发出的输出会再发一遍。
  • exactly-once(端到端):Flink checkpoint(内部状态一致)+ Kafka 事务 Sink(两阶段提交,2PC)配合:预提交的输出在 checkpoint 完成前对下游不可见,checkpoint 成功才 commit。代价是输出可见性延迟被拉长到「约一个 checkpoint 间隔」,且要求下游是事务性 Sink——向量库通常不是,这正是思考题 1 的切入点。

工程共识:当 Sink 不支持事务时,退而求其次的黄金组合是 at-least-once + 幂等写入(idempotent upsert)——效果上等价于 exactly-once(业界称 effectively-once),实现成本低一个数量级。

2.2 Embedding 管道三模式

1. Batch 全量

  • 输入:全量文档 Parquet
  • 输出:向量库 collection docs_v3
  • 适用:冷启动、embedding 模型从 v2→v3

2. Streaming 增量

  • 触发:doc.updated 事件
  • 流程:分块 → embed batch(32) → upsert by doc_id
  • 适用:日常更新

3. Backfill 补历史

  • 触发:新字段需要 re-embed(如加了 title prefix)
  • 流程:扫描 updated_at < T 的 ID 列表,限速回填,避免打满 GPU

2.3 Embedding 版本切换 SOP

步骤动作
1部署新 embedding 模型 v2.1 到独立 worker pool
2双写:新文档同时写 index_v2_0index_v2_1(或 shadow index)
3Backfill 历史文档到 index_v2_1,监控 lag
4RAG 路由切读 v2_1,保留 v2_0 24–72h 回滚
5下线旧 index,更新 registry metadata

这与 L5.2 模型发布的 blue-green 同源——索引也是模型产物

2.4 背压(Backpressure):流式管道的守恒定律

流式系统只有一条铁律:长期看,消费吞吐必须 ≥ 生产速率,否则队列无界增长。设生产速率 λ(docs/s)、消费吞吐 μ(docs/s),则积压增长速度 = λ − μ;只要 λ > μ 持续 T 秒,积压量就是 (λ − μ) × T,端到端延迟随之线性恶化。

Embedding 管道的特殊性在于 μ 由 GPU 推理决定μ = batch_size / 单批耗时。它对三件事极其敏感——模型换大一号(单批耗时翻倍 → μ 腰斩)、batch 攒不满(低峰期反而每条更慢)、backfill 抢占同一 GPU pool。健康的管道要有三层防线:

  1. 可观测:消费组 lag 及其斜率(lag 只涨不落 = 结构性失衡,比绝对值更早报警);
  2. 可缓冲:Kafka 本身就是背压缓冲区,但缓冲只对尖峰有效,对稳态失衡无效;
  3. 可伸缩/可降级:按 lag 自动扩 embedding worker,或降级策略(低优先级 topic 延后、动态加大 batch)。

思考题 2 会用真实数字把这条守恒定律算一遍。

动手实践:最小流式 Embedding 管道

实验目标:不装任何中间件,用纯 Python 的 queue + threading 复刻「producer → 攒批 → embed → 写库」的最小流式管道,亲手观察攒批(batching)与消费吞吐的关系;再用 Docker Compose 起一个 KRaft 单节点 Kafka 做真实演练入口。

A. Mini 流式管道(纯 Python,无 Kafka)

# mini_stream_embed.py
import queue, threading, time, hashlib

docs_q = queue.Queue()
results = []

def producer():
for i in range(20):
docs_q.put({"doc_id": f"d{i}", "text": f"content {i}"})
time.sleep(0.05)
docs_q.put(None)

def embed_batch(batch):
return [{"doc_id": d["doc_id"], "vec": hashlib.md5(d["text"].encode()).hexdigest()[:8]} for d in batch]

def consumer():
batch, BATCH = [], 4
while True:
item = docs_q.get()
if item is None:
if batch:
results.extend(embed_batch(batch))
break
batch.append(item)
if len(batch) >= BATCH:
results.extend(embed_batch(batch))
batch = []

t0 = threading.Thread(target=producer)
t1 = threading.Thread(target=consumer)
t0.start(); t1.start(); t0.join(); t1.join()
print(f"embedded={len(results)} sample={results[:2]}")

B. Docker Compose:Kafka 单节点(KRaft 模式,本地演练用官方 apache/kafka 镜像)

# docker-compose-stream.yaml(节选)——KRaft 单节点:broker + controller 一体,无 ZooKeeper
# 镜像 tag 按 Kafka 4.x 线的当前 release 替换即可,配置字段同构
services:
kafka:
image: apache/kafka:4.0.0
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
ports: ["9092:9092"]
docker compose -f docker-compose-stream.yaml up -d
# 生产环境用 kafka-console-producer 写入 doc.updated 事件
# Flink SQL: CREATE TABLE ... WITH ('connector'='kafka', ...)

C. Backfill 限速脚本示意

# backfill.py — 限速 100 doc/s,避免 GPU 打满
import time
for doc_id in open("ids_to_backfill.txt"):
embed_and_upsert(doc_id.strip())
time.sleep(0.01)

踩坑预警

  1. 流式 job 无窗口groupBy 无界 state 最终 OOM。
  2. embedding 与 chunk 策略变更未 bump index version:检索质量静默下降(见 L5.7 事故案例)。
  3. backfill 不限速:挤占在线 streaming GPU,P99 爆炸。
  4. Kafka 分区数 < 并行度:Flink 并行 task 空闲。

配套代码ai-infra-labs mini_stream_embed.py

深入思考

下面三题每题先给题干,再用 <details> 折叠一份图文并茂的参考答案。建议先合上答案自己想 3 分钟,再展开对照。

思考题 1:exactly-once 太贵,at-least-once 会重——向量库怎么办?

2.1 节说端到端 exactly-once 需要「Kafka 事务 + 两阶段提交 Sink」,但 Milvus/pgvector 这类向量库 Sink 通常不支持分布式事务。假设你的 Flink job 以 at-least-once 运行,某晚 OOM 重启后重放了 checkpoint 之后的 8000 条 doc.updated 事件(即「背景与现状」的症状 2)。请回答:重复写入向量库会造成什么检索层面的后果?为什么「按 doc_id upsert」还不够、chunk 级主键必须是确定性的?给出一套让重放无害的幂等(idempotent)设计。

展开参考答案(含幂等写入设计图 + 算一遍)

结论:向量库拿不到事务性 Sink 时,正确解法不是硬上 exactly-once,而是 2.1 节说的 at-least-once + 幂等 upsert(effectively-once)。关键在主键设计——主键必须由输入内容确定性导出(如 doc_id + chunk_seq + chunk_hash),使「同一条事件写两遍」落到同一主键上被覆盖而非新增。若主键含任何非确定成分(自增 ID、UUID、处理时间戳),重放就会制造新行,重复向量挤占 top-k 名额、检索质量静默下降且无任何报错。

为什么只按 doc_id upsert 还不够:一个 doc 会切成多个 chunk。若文档从 12 个 chunk 改到 10 个,按 doc_id + chunk_seq upsert 只覆盖前 10 个,旧的第 11、12 号 chunk 变成幽灵残留。所以完整幂等还需「先按 doc_id 删旧 chunk 集、再写新集」或对旧 chunk 打 tombstone——这正是 2.2 节「upsert by doc_id」在 chunk 粒度上的严格化。

用具体数字算一遍:向量库共 100,000 个 chunk,重放 8000 条事件、平均每 doc 6 个 chunk,非幂等 insert 会新增约 8000 × 6 = 48000 个重复 chunk(库膨胀 48%)。设 5% 的查询其 top-10 命中了被重放的热点文档,每次命中平均被 23 个重复 chunk 挤掉 12 个本应出现的多样性结果——检索 recall 掉了,但所有基础设施监控全绿。这就是「症状 2」难查的原因:错误不在延迟或错误率,而在结果分布里。

方案不丢不重Sink 要求成本
at-most-once最低但不可用
at-least-once + 随机主键 insert低,但重放有害
exactly-once(2PC)事务性 Sink高,向量库通常不支持
at-least-once + 确定性主键 upsert✅(效果上)支持 upsert 即可低,推荐

回看 2.1 节的「工程共识」:这套设计把「不重」的责任从传输层(事务)挪到了存储层(幂等键),是流式 Embedding 管道的标准答案。

思考题 2:背压算一遍——模型升级后 lag 为什么只涨不落?

复现「背景与现状」的症状 1:embedding worker 以 batch=32、单批 50 ms 的速度消费;某晚模型从 base 升级为 large,单批耗时变为 100 ms。日间稳态生产速率约 400 docs/s、早高峰 2 小时冲到 800 docs/s。结合 2.4 节的守恒定律:① 分别算出升级前后的消费吞吐 μ;② 说明为什么升级前系统「扛得住高峰」、升级后 lag 却结构性上涨;③ 给出至少两种止血手段并算出各自的 lag 消化时间。

展开参考答案(含背压水位图 + 算一遍)

结论:升级前 μ = 32/0.05 = 640 docs/s,高于稳态 400 也只是短暂低于高峰 800,Kafka 缓冲吃下尖峰后能回落;升级后 μ = 32/0.10 = 320 docs/s,低于稳态 400——这不再是「尖峰」而是「结构性失衡」,lag 以 80 docs/s 恒定速度上涨,任何缓冲都救不了,只能改 μ(扩 worker/加大 batch)或改 λ(降级限流)。判断标准就一条:λ 与 μ 谁大——2.4 节三层防线里「lag 斜率」比绝对值更早暴露这一点。

逐步算一遍(下列数字已用 Python 逐项验算):

  1. 升级前 μ = 32 / 0.05 = 640 docs/s。稳态 λ=400 < 640,健康;早高峰 λ=800 > 640,2 小时积压 (800−640) × 7200 = 1,152,000 条,但高峰一过,以 640 − 400 = 240 docs/s 消化,1,152,000 / 240 = 4800 s = 80 min 清零——lag 曲线是「山峰」形,涨得上去落得下来
  2. 升级后 μ = 32 / 0.10 = 320 docs/s < λ=400。缺口 80 docs/s 与流量高低无关地持续存在:2 小时积压 80 × 7200 = 576,000 条;即使生产完全停止,320 docs/s 也要 576,000 / 320 = 1800 s = 30 min 才能清完——lag 曲线是「爬坡」形,只涨不落,这正是症状 1 里「斜率为正」的报警意义。
  3. 止血手段 A——扩 worker:GPU worker 从 1 个扩到 2 个,μ = 640;消化 576,000 条积压需 576,000 / (640 − 400) = 2400 s = 40 min
  4. 止血手段 B——加大 batch:GPU 推理近似「单批耗时随 batch 亚线性增长」,若 batch 32→64 单批耗时仅从 100 ms 涨到 160 ms,则 μ = 64 / 0.16 = 400 docs/s——刚好持平 λ,能止涨但消化不了存量,需与手段 A 组合。

这题的真正教训:embedding 模型升级是「容量变更」而不只是「模型变更」——上线前必须重测单批耗时、按 2.4 节守恒定律重新核算 μ 与峰值 λ 的余量,把「容量核算」写进 2.3 节版本切换 SOP 的第 1 步。

思考题 3:模型 v2→v3 切换,双写窗口里的「新文档」怎么保证不丢不错?

2.3 节 SOP 做 embedding 模型 v2→v3 切换:backfill 存量 1000 万 chunk 需要十几个小时,期间流式管道还在源源不断产生新文档。请回答:① 为什么必须双写(只写新索引会怎样)?② backfill 与流式双写同时进行时,同一个 doc 可能被两条路径先后写入 index_v3,如何避免「backfill 的旧内容覆盖流式的新内容」?③ 用 200 docs/s 的 backfill 限速与单条 0.02 GPU·s 的推理成本,算一算这次切换的时间与 GPU 账单。

展开参考答案(含双写/回填时序图 + 算一遍)

结论:双写的本质是让 index_v3 在「追平存量」的十几个小时里不漏掉任何增量——若只写 v3 而 RAG 还在读 v2,切换前的新文档在 v2 里缺失、用户立刻搜不到新内容;若只写 v2,v3 上线时又缺了这段窗口。而 backfill 与流式的写入竞争,用「版本/时间戳条件写」解决:每条向量带 updated_at(或单调版本号),backfill 采用 upsert-if-older——仅当库内该主键的 updated_at 更旧时才写入,天然让「流式赢过回填」。检索结果永不混用两个模型的向量空间:路由以 index 为单位原子切换,这是 2.3 节 SOP 第 4 步的深层原因。

为什么不能「新文档只写 v3、提前切读」:v2 与 v3 的向量不在同一语义空间,一条查询向量(由 v3 编码)去检索 v2 向量,相似度分数没有可比性——所以既不能混索引查询,也不能在 v3 未追平前切读。双写 + 整索引原子切换是唯一干净解。

用具体数字算一遍(已用 Python 验算):

  1. backfill 时长:1000 万 chunk ÷ 200 docs/s(限速,见 2.2 节模式 3)= 50,000 s ≈ 13.9 h——这就是双写窗口的最小长度,说明「切换是一个以天计的流程,不是一次发布」。
  2. GPU 账单10,000,000 × 0.02 GPU·s = 200,000 GPU·s ≈ 55.6 GPU·h;按每卡每小时 2 美元计约 111 美元——真正贵的不是 GPU,而是双写期间双倍的向量库写入与存储(两份索引并存 1~3 天)。
  3. 竞争窗口有多大:13.9 h 内以稳态 400 docs/s 计,约 2000 万条流式写入与 backfill 并发落到 v3;若不做 upsert-if-older,任何在窗口内被编辑过的 doc 都有概率被回填的旧快照回滚——按 1% 的文档窗口内被编辑估算,就是约 10 万个 chunk 的静默数据回退。

收尾三件事(对应 2.3 节 SOP 第 4~5 步):切读前对账(v2/v3 条数差 < 0.1%)+ 抽样 A/B 检索质量;保留 v2 至少 24 h 应对「v3 检索质量回退」类回滚;下线时同步更新 registry,让「哪个索引由哪个模型版本生成」永远可追溯——索引也是模型产物

延伸阅读

1. 官方文档(概念必读)

  • Kafka 官方文档 §Design(投递语义、事务与 EOS)与 §KRaft — 理解 4.x 线「无 ZooKeeper」架构与 exactly-once 语义边界。
  • Flink 官方文档 §Stateful Stream Processing(checkpoint、两阶段提交 Sink)与 §Event Time / Watermark — 本文 2.1 节的完整展开。
  • Streaming Systems(Tyler Akidau 等)— 流式语义的系统性经典,「What / Where / When / How」四问框架。

2. 相关高 Star 仓库

  • apache/kafka / apache/flink — 直接读 EOS 与 checkpoint 相关的设计文档与源码注释。
  • milvus-io/milvus — 看其 upsert 与主键去重语义,验证思考题 1 的幂等设计前提。

3. 站内衔接