数据截至 (上游 commit a51131d3fc7a)
摄取管线与高级 agent
30 秒导读: 前面几章讲的都是「用户来了一个问题、agent 怎么答」。这一章讲两件更靠前、也更复杂的事: ①知识是怎么进库的——一份 PDF / 一个 GitHub 仓库 / 一个 Confluence 空间,如何被解析、切块、向量化,最终变成第 4 章 检索层 能查到的东西;②两种比普通工具循环更重的 agent——
ResearchAgent(自己拆问题、分步调研、写带引用的报告)和WorkflowAgent(执行用户画好的节点图)。
本章在全书中的位置:它是 01 从请求到 agent→02 工具循环→03 工具体系→04 检索层→05 LLM 抽象 之后的「上游 + 加长版」。检索能查到东西的前提,是本章前半段先把东西写进去。
1. 这是什么(零基础也能懂)
先分清本章讲的两块内容,它们方向相反:
| 这块 | 方向 | 一句话 | 谁是它的下游/服务对象 |
|---|---|---|---|
| 摄取管线(ingestion) | 写入 | 把原始文档变成可检索的向量(和可选的知识图谱) | 第 4 章的检索器从这里读 |
| 高级 agent | 读取/编排 | 比普通 agent 多几层编排:多步调研、节点图执行 | 直接面向用户,产出报告/工作流结果 |
摄取:一句话直觉
把它想成一条流水线的传送带:一头进原料(各种格式的文件、网页、云盘文档),中间几道工序(解析成纯文本 → 切成小块 → 每块算成一个向量),另一头把成品码进货架(向量库)。整条带子跑在后台异步 worker(Celery)上,因为解析一个大 PDF、给一万个块算 embedding 都很慢,不能卡住用户的 HTTP 请求。
高级 agent:一句话直觉
ResearchAgent像一个做课题的研究员:拿到一个大问题,先判断要不要追问澄清,再把它拆成几个子问题,一个个去检索/搜资料,最后把发现汇总成一篇带编号引用的报告。WorkflowAgent像一条用户自己画的装配线:用户在画布上拖出「节点 + 连线」的图(agent 节点、代码节点、条件分支节点……),这个 agent 就照着图一个节点一个节点地跑。
2. 顶层全景(摄取管线怎么转)
2.1 一张图:从原料到货架
┌─────────────── 数据来源(三类)────────────────┐
本地文件 ─────► │ file/* 远程 URL ─► remote/* 云盘/wiki ─► connectors/* │
(上传/zip) └──────────────────────────┬───────────────────────────────┘
│ 下载 / 抓取到临时目录
▼
① 解析 SimpleDirectoryReader + 各格式 Parser
(PDF/docx/pptx/html/md/audio/表格… → 纯文本 Document)
│
▼
② 切块 ChunkerCreator → 5 种策略之一
(按 token 窗口 / 递归 / markdown / 父子 / 语义)
│
▼
③ 向量化 embed_and_store_documents
(逐块 embed → 写向量库;ingest_chunk_progress 记进度)
│
┌────────────────────────┴───────────────────────┐
▼(ClassicRAG 立即可用) ▼(仅 graphrag 源,异步分支)
第 4 章的检索器可查 extract_graph_for_source
(每块再过一次 LLM 抽实体/关系 → 图谱表)
怎么读这张图: 从上到下是同一条传送带,①②③ 顺序执行;右下角的图谱抽取是一条可选的旁路——只有 graphrag 类型的源才会在向量化完成后再异步跑一遍。
2.2 部件与职责
| 部件 | 干什么 | 在哪(相对克隆根) |
|---|---|---|
SimpleDirectoryReader | 遍历目录、按扩展名挑解析器、产出 Document 列表 | application/parser/file/bulk.py:157 |
| 各格式 Parser | 把一种格式解成文本(Docling 为主) | application/parser/file/* |
| 远程加载器 | 抓 URL / sitemap / GitHub / Reddit / S3 | application/parser/remote/* |
| 连接器 | 拉 Google Drive / SharePoint / Confluence | application/parser/connectors/* |
Chunker 系列 | 把长文本切成 token 受限的块 | application/parser/chunking.py、chunking_strategies.py |
embed_and_store_documents | 逐块算 embedding、写向量库、记断点 | application/parser/embedding_pipeline.py:149 |
ingest_chunk_progress 表 | 记「已 embed 到第几块」,供断点续跑 | 由 IngestChunkProgressRepository 读写 |
ingest_worker 等 Celery 任务 | 串起①②③、发进度事件、收尾触发图谱 | application/worker.py:490 |
extract_graph_for_source | 每块过 LLM 抽实体/关系,建图谱 | application/graphrag/extraction.py:149 |
GraphStore | 图谱的 pgvector 存储层 | application/graphrag/store.py:76 |
2.3 主线走一遍(高层,不进代码)
一次本地文件上传:HTTP 路由把文件落到存储、ingest_worker 入队 → worker 把文件拉到临时目录、SimpleDirectoryReader.load_data 解析成文本 → ChunkerCreator 按源配置切块 → embed_and_store_documents 逐块向量化并落库、ingest_chunk_progress 每块打一个勾 → 发 source.ingest.completed 事件 → 若是 graphrag 源,_maybe_enqueue_graph_extraction 再排一个图谱任务(application/worker.py:729)。
3. 摄取:逐层拆开
3.1 解析层:各格式各自为政
它要解决的小问题: PDF、Word、音频、Excel……字节结构天差地别,但下游只想要「纯文本」。解析层就是把 N 种格式收敛成一种 Document。
入口是 SimpleDirectoryReader(application/parser/file/bulk.py:157):它遍历目录,对每个文件按扩展名从一张扩展名 → 解析器实例的字典里取解析器。这张字典由 get_default_file_extractor(同文件 :28)构建,默认走 Docling(一个文档转换库),只有少数格式保留了专用解析器:
| 扩展名 | 默认解析器 | 说明 |
|---|---|---|
.pdf | DoclingPDFParser | 可选 OCR(ocr_enabled) |
.docx / .pptx / .xlsx | DoclingDocxParser / DoclingPPTXParser / DoclingXLSXParser | Office 三件套 |
.html / .xhtml / .xml | DoclingHTMLParser / DoclingXMLParser | |
.csv | DoclingCSVParser | 另有 PandasCSVParser/ExcelParser 走 fast 引擎 |
.md / .mdx | MarkdownParser(专用,非 Docling) | 保留专门处理 |
.rst | RstParser | |
.json | JSONParser(专用) | |
.png / .jpg / .tiff / .webp… | DoclingImageParser(开 OCR)或 ImageParser | 图像 |
.epub | EpubParser | |
音频(见 SUPPORTED_AUDIO_EXTENSIONS) | AudioParser | 语音转写 |
可摄取的扩展名白名单集中在 application/parser/file/constants.py:文档类 SUPPORTED_SOURCE_DOCUMENT_EXTENSIONS、图像类、音频类合并成 SUPPORTED_SOURCE_EXTENSIONS:23——这张白名单是「不可信文档」的第一道闸(不在名单里的扩展名直接拒)。
另一个解析入口:document_reader.py。 这个文件不是给批量入库用的,而是给 第 3 章 的 read_document 工具用的——让 agent 在对话中即时解析一份文档。它复用同一套 Docling 解析器,但额外套了一层面向不可信输入的护栏,值得单独看:
| 护栏 | 做什么 | 代码 |
|---|---|---|
| 扩展名 白名单 | 只认 SUPPORTED_SOURCE_EXTENSIONS | parse_document_bytes application/parser/document_reader.py:434 |
| 字节上限 | 超 DOCUMENT_PARSE_MAX_BYTES / 沙箱上限直接拒 | 同文件 :399 |
| zip 炸弹检测 | docx/xlsx/pptx/epub 本质是 zip,先查条目数与解压后总大小 | _reject_zip_bomb:108 |
safe_filename 落盘 | 恶意文件名被消毒后才落临时文件,用完即删 | :393、finally :415 |
| 头+尾窗口截断 | LLM 看到的视图截到字节预算内,但完整解析仍持久化 | truncate_text_head_tail:46、bound_parse_payload:62 |
这里有个值得学的取舍:解析产出的是全文,截断只发生在「给 LLM 看的视图」上(_bounded:502 明确「永不在此截断」,注释说明完整文本进 data artifact,视图在 bound_parse_payload 才收窄)。这样既不让大文档撑爆上下文,又不丢原文。
远程与连接器是解析层的另外两个入口,结构上都是工厂 + 一组加载器:
- 远程
RemoteCreator(application/parser/remote/remote_creator.py:11),按type取加载器:url→WebLoader、sitemap→SitemapLoader、crawler→CrawlerLoader、reddit→RedditPostsLoaderRemote、github→GitHubLoader、s3→S3Loader(另有crawler_markdown、telegram等文件)。 - 连接器
ConnectorCreator(application/parser/connectors/connector_creator.py:9):confluence/google_drive/share_point,且每个连接器配一套 OAuth 认证 provider(create_auth:49)——因为云盘/wiki 要授权才能读。
远程/连接器加载完文件后,汇入同一条 ②③ 传送带(worker 里 remote_worker、ingest_connector 也是 ChunkerCreator + embed_and_store_documents)。
3.2 切块层:5 种策略,一套 token 预算
它要解决的小问题: embedding 模型和上下文都有长度上限,一篇长文必须切成小块,且切法直接影响检索质量。
派发在 ChunkerCreator;每种策略是一个注册进去的实现。经典策略 Chunker(application/parser/chunking.py:11,注册名 classic_chunk)是按 token 窗口切:
# 示意,非源码。经典切块的核心思路
header, body = separate_header_and_body(text) # 前 3 行当"表头",可复制到每块
while 还有 body 没切完:
end = 当前位置 + max_tokens - len(header_tokens)
chunk = header_tokens + body[当前位置:end] # 第 0 块或 duplicate_headers 时带表头
emit(chunk); 当前位置 = end # 重点看:表头只有第 0 块默认带
真实实现见 split_document:44 与 classic_chunk:69:小于 min_tokens 的文档原样保留不再切,在 [min,max] 区间的直接过,超 max_tokens 的才 split_document。默认 max_tokens=2000 / min_tokens=150。
另外四种策略在 chunking_strategies.py,共享 _BaseStrategyChunker 的 token 工具(同一套 tiktoken 编码,预算口径一致):
| 策略(注册名) | 切法 | 亮点 | 类 |
|---|---|---|---|
recursive | 按 \n\n → \n → ". " 分隔符层层降级,最后硬切 token | 小碎片再 _merge_to_min 合并到过 min_tokens | RecursiveChunker:90 |
markdown | 按 ^#{1,6}\s 标题切段,超长段再 token 切 | 尊重文档结构 | MarkdownChunker:131 |
parent_child | 先切大「父窗口」,再切小「子块」;子块入向量、父文本进 extra_info["parent_text"] | 检索命中小块、可回放大块上下文 | ParentChildChunker:172 |
semantic | 句子逐个 embed,在余弦距离**高百分位(95)**处断开 | 边界落在话题切换处;失败自动降级到 recursive | SemanticChunker:220 |
semantic 的巧妙点:它把整篇的句子一次批量 embed,算相邻句的余弦距离,取 95 百分位当阈值断句(_breakpoints:246);任何异常(句子太少、embedding 失败、距离退化)都 _fallback 到递归策略(:237),保证入库永不崩。
一个重要边界:切块策略是入库时才决定的;换策略必须重新入库(源码注释里标为 D8,见
chunking_strategies.py:6)。
3.3 向量化层:逐块 embed + 可续跑的断点
它要解决的小问题: 给一万个块算 embedding 可能跑几分钟,中途一次限流/网络抖动不该让整批白干。
核心是 embed_and_store_documents(application/parser/embedding_pipeline.py:149)。它不是「一把梭」,而是逐块循环 + 每块打勾:
┌── ────── ingest_chunk_progress(每个 source 一行)────────┐
attempt_id 匹配? ──►│ 匹配(同一任务的 Celery 重试)→ 从 last_index+1 续跑 │
│ 不匹配(全新 sync/reingest)→ 重置 checkpoint,从 0 重建 │
└──────────────────────────────────────────────────────┘
for idx in [loop_start, total):
add_text_to_store_with_retry(store, doc) # 单块内 @retry(3,5,backoff2)
_record_progress(idx) # 落 checkpoint(尽力而为)
throttle 每涨 1% 发一次 SSE(upload toast)
出错 → 保存已成向量 + break + 抛 EmbeddingPipelineError(让 Celery autoretry)
几个设计要点,都能直接借鉴:
- 断点续跑靠
attempt_id语义(_init_progress_and_resume_index:68):同一任务的 Celery 自动重试携带同一个attempt_id,于是从last_index+1接着跑;而全新的 sync/reingest 是不同attempt_id,会重置 checkpoint 从 0 重建——这条正是「防止一个已完 成的旧 checkpoint 悄悄让下一次 sync 变成 no-op」。 - 失败必须重新抛出(
:335):即便部分成功也raise EmbeddingPipelineError。注释讲得很清楚——如果吞掉异常,任务体会返回成功,with_idempotency就会把一个部分索引标成completed并缓存 24h,毒化后续。 - 防御性 tripwire
assert_index_complete(:117):worker 在 embed 之后再查一次ingest_chunk_progress,embedded < total就抛错,防止任何未来的「吞异常」路径把残缺索引当完整缓存。 - 心跳线程(
worker.py:_start_ingest_heartbeat:164):后台每 30s 更新ingest_chunk_progress.last_updated,让长时间跑的 embed 不被判定为「卡死」。 - FAISS vs 其它库的分叉:FAISS 需至少一个文档才能建索引,所以用
docs[0]播种、从索引 1 起循环(:226);续跑时则从磁盘加载已有 FAISS 索引再追加(:213)。
3.4 Celery 编排:worker.py 把一切串起来
摄取的所有慢活都跑在 application/worker.py 的 Celery 任务里,主要几个:
| 任务 | 触发 | 干什么 |
|---|---|---|
ingest_worker:498 | 本地文件/zip 上传 | 解析→切块→向量化→发事件→触发图谱 |
reingest_source_worker:762 |