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
Documents SentenceSplitter切块 TitleExtractor抽标题→metadata Embedding向量化 就绪 Node 每个方块都是 TransformComponent(输入 Node → 输出 Node),像乐高一样串起来
切块器、元数据抽取器、嵌入模型都是"输入 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 份文档,你只改了 财务制度.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(D3) 切块(D4) 嵌入(D6) 向量库(D9) ✅ 可检索 第 1 周闭环:原始文档 → 一条管道 → 直接进向量库、随时可查
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 从"嵌入是什么"讲起。
← Day 04 Day 06 · 嵌入 →