Day 05 / 共 20 天 · 第 1 周 RAG 全景与数据摄入
摄入管道 Ingestion(把①②③串起来)
第 1 周收官。前几天的"读→切→嵌"其实可以打包成一条流水线,一键跑完、还能缓存。今天看 IngestionPipeline 怎么串联步骤、怎么用缓存省下重复处理(尤其省嵌入 API 的钱)、怎么增量更新。
📍 你在整条链的位置
① 读+
② 切+
③ 嵌=
Ingestion 管道(打包①②③)→
④ 存索引
L01
为什么要"管道"
🤔 每次建库都手写"读→切→抽元数据→嵌入",烦不烦?
更要命的是:你有 1000 份文档已建好库,现在只改了 1 份。重新跑一遍全部?——切块、嵌入 1000 份,白花 999 份的时间和嵌入 API 的钱(嵌入是按 token 收费的!)。
💡 IngestionPipeline = 可复用 + 可缓存 + 可增量 的摄入流水线
把步骤列成一条管道,一键跑完;缓存让没变的文档直接用旧结果(不重嵌);增量让只处理变化的那 1 份。把摄入这件事"工程化"。
L02
转换(Transformation)是乐高块
💡 关键抽象:TransformComponent(schema.py:191)
所有处理步骤都实现同一个接口:输入一批 Node、输出一批 Node。
__call__(nodes) -> nodes。切块器、元数据抽取器、嵌入模型都是"输入 Node、输出 Node"的转换——所以能任意串成管道。
"统一接口"的威力(本系列反复出现)
因为切块/抽元数据/嵌入都长一个样(输入 Node、输出 Node),它们就能像乐高一样随意拼接。你甚至能加自定义转换(比如"过滤掉太短的块""给每块加公司名前缀")插进管道任意位置。
L03
run_transformations:就是个 for 循环
ingestion/pipeline.py:72-113——串联逻辑简单到出乎意料:
def run_transformations(nodes, transformations, ...):
for transform in transformations: # 依次跑每个转换
nodes = transform(nodes, **kwargs) # 上一个的输出 = 下一个的输入
return nodes
💡 简单是因为接口统一
正因为每个转换都是"输入 Node、输出 Node",串联就是"把上一个的输出喂给下一个"——一个 for 循环搞定。好设计让核心代码极简。有
arun_transformations(:115)异步版,嵌入调 API 时并发加速。L04
IngestionPipeline:打包好用
📝 显式用管道(能精细控制每一步)
from llama_index.core.ingestion import IngestionPipeline
pipeline = IngestionPipeline(transformations=[
SentenceSplitter(chunk_size=512), # 切块(Day4)
TitleExtractor(), # 抽标题当 metadata
OpenAIEmbedding(), # 嵌入(Day6)
])
nodes = pipeline.run(documents=documents) # 一键跑完,得到就绪 Node读法:
pipeline.run 内部就是 run_transformations + 缓存逻辑。Day 01"5 行 RAG"里的 VectorStoreIndex.from_documents 内部就在做类似的隐式摄入。想省事用 from_documents,想精细控制(自定义步骤、复用缓存)就显式用 Pipeline。L05
缓存:只改 1 份不重嵌 999 份
💡 缓存键 = 输入内容的 hash + 转换的配置
每个转换执行前,先算"这批 Node 的 hash + 这个转换"的键,查
IngestionCache——命中就直接用结果,不重跑。📝 第二次跑管道时发生了什么
1000 份文档,你只改了
• 999 份没变的 → hash 一样 → 缓存命中 → 切块/嵌入结果直接取缓存,0 次 API 调用
• 1 份改了的 → hash 变了 → 缓存未命中 → 重新切块 + 嵌入
→ 嵌入 API 调用从 1000 次降到 1 次,省钱 999 倍。
财务制度.pdf:• 999 份没变的 → hash 一样 → 缓存命中 → 切块/嵌入结果直接取缓存,0 次 API 调用
• 1 份改了的 → hash 变了 → 缓存未命中 → 重新切块 + 嵌入
→ 嵌入 API 调用从 1000 次降到 1 次,省钱 999 倍。
"内容 hash 做缓存键"是通用手法
回想你学过的 APISIX(conf_version 版本号缓存)、wasm-go(规则备份)——都是"用内容/版本做键,变了才失效"。Node 的
hash 字段(Day 02,基于内容算)正是为此。嵌入是 RAG 最花钱的部分,这个缓存价值巨大。L06
增量更新文档(UPSERTS)
🤔 知识库不是建一次就不动的
文档会新增、修改、删除。每次全量重建太浪费——怎么"只处理变化的"?
💡 配 docstore + UPSERTS 策略
Pipeline 接一个
docstore 后,按 Document 的 id_ + hash 判断:新文档 → 处理入库;改过的(hash 变)→ 重新处理并更新;没变的 → 跳过;删掉的 → 从库里移除。读法:回收 Day 03 的
filename_as_id——用稳定的 id 才能正确判断"是不是同一个文档"(不然改了名就当成新文档,重复入库)。增量更新是生产级 RAG 维护知识库的必备能力。L07
直连向量库:第 1 周闭环
💡 Pipeline 可以直接把结果写进向量库
传
vector_store 参数(Day 09),跑完管道后嵌入好的 Node 直接进向量库,随时可检索——不用再手动建索引。Reader → 切块 → 嵌入 → 向量库,一条 Pipeline 打通"离线建库"全流程。
L08
今日小结 + 动手(第 1 周收官)
🧠 第 1 周你应该能回答
- 为什么需要摄入管道?解决"重复处理"和"增量更新"什么问题?
- TransformComponent 为什么能让切块/嵌入像乐高一样拼接?
- 缓存怎么用内容 hash 做键?为什么能省嵌入 API 的钱?
- UPSERTS 怎么实现"只处理变化的文档"?
✋ 动手
cd /Users/bitmart/work/codes/github/llama_index/llama-index-core/llama_index/core
sed -n '72,140p' ingestion/pipeline.py
grep -n "class IngestionPipeline\|class IngestionCache\|docstore_strategy\|def run" ingestion/pipeline.py | head
sed -n '191,207p' schema.py # TransformComponent
下周预告 · 第 2 周:进入 RAG 最神奇的部分——嵌入(文字怎么变成"能算相似度"的向量)、索引(怎么组织成可检索的结构)、存储(数据存哪、怎么持久化)。Day 06 从"嵌入是什么"讲起。