L2.9 流式 ETL 与 Embedding 管道
三维坐标
layer: L2(数据与算子层)|level: Senior|pillar: 数据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 为准。