数据截至 (上游 commit 8c51b8dc5408)
摄取流水线:一个文件怎么变成可被搜索的东西
本章讲什么: 从
POST /ingest/file到文档status=completed,中间发生了什么、在哪一步会卡住、解析失败时它怎么补救。
1. 为什么要拆成同步 + 异步两段
一份 300 页的 PDF,ColPali 路径下要渲染 300 张图、跑 300 次视觉模型前向。这个时间量级是分钟,不可能压在一个 HTTP 请求里。
所以接口层只做三件轻活,做完立刻返回:
POST /ingest/file
|
|-- ① Postgres 插一行 Document,status = processing
|-- ② 原始字节上传到对象存储
|-- ③ Redis 塞一个 arq 任务
|
v
返回 Document(调用方拿到 external_id,之后自己轮询 status)
对应 core/services/ingestion_service.py:552 的 ingest_file_content:第 610 行 store_document 建行,第 622 行 _upload_content_bytes 上传,第 686 行 redis.enqueue_job("process_ingestion_job", **job_payload) 入队。
顺序是刻意的:先建库行再上传。任何一步失败都会走 _mark_document_failed(core/services/ingestion_service.py:399)把库里那行标成 failed,而不是留下一个孤儿。
任务 ID 就是去重键
入队的 payload 里有一个字段值得单独看(core/services/ingestion_service.py:432-434):
"_job_id": f"ingest:{document_id}",
"_expires": timedelta(days=7),
arq 用 _job_id 做去重。同一个 document_id 重复排队时,第二次 enqueue_job 返回 None,代码把这当成正常情况打一条 info 日志(core/services/ingestion_service.py:687-688)。重复点"重新摄取"不会真的跑两遍。
2. worker 的六步流水线
主函数是 core/workers/ingestion_worker.py:448 的 process_ingestion_job。它在第 553 行把 total_steps 定死成 6,每一步都往数据库写一次进度,前端因此能画进度条。
[1] Downloading file 从对象存储把原始字节拉回来
|
v
[2] Parsing file 决定走"文本路"还是"图片路"
|
v
[3] Splitting into chunks 文本路才切块;图片路直接一页一块
|
v
[4] Generating embeddings 文本路走普通嵌入,图片路走 ColPali
|
v
[5] Storing chunks 按批写向量库 + 内容写对象存储
|
v
[6] Finalizing status = completed,清掉 progress
进度写入靠 update_document_progress(core/workers/ingestion_worker.py:170),它优先调用 update_document_system_metadata_fields 做服务端字段合并,避免先读一整行 system_metadata(那里面装着全文)再写回去。
一个容易忽略的分叉:要不要解析文本
第 646-649 行是全流程最关键的分支:
colpali_native_format = is_colpali_native_format(mime_type)
skip_text_parsing = using_colpali and colpali_native_format
is_colpali_native_format(core/storage/utils_file_extensions.py:75)对这些类型返回真:
| 类型 | MIME |
|---|---|
| 任意图片 | image/* |
application/pdf | |
| 医学影像 | application/dicom |
| Word | ...wordprocessingml.document、application/msword |
| PowerPoint | ...presentationml.presentation 等三种 |
只要开了 ColPali 且格式在这张表里,文本解析整个跳过。 Docling 不跑、OCR 不跑、切块不跑。这是性能上最大的一笔节省,也是"页即块"能成立的前提。
一个小惊喜:HTML 先转 PDF
第 619-632 行,遇到 text/html 会先用 WeasyPrint 渲染成 PDF,再当 PDF 处理。思路是"模拟打印出来的样子",让网页也能享受版面感知的检索。转换失败就退回按原始 HTML 处理,不中断。
3. 文本路:解析器 的多条分支
走文本路时,MorphikParser.parse_file_to_text(core/parser/morphik_parser.py:617)按文件类型分流。这不是一条 if-else 链,而是几条独立的快慢通道:
| 输入 | 走哪条路 | 为什么 |
|---|---|---|
.txt/.md/.json/.csv/.yaml 等 | 直接 decode,不解析 | 没必要动重型解析器(_is_plain_text_file:273) |
.xlsx/.xlsm | openpyxl 转 markdown 表格 | 注释原文说比 Docling 快"亚秒 vs 分钟",且标签和数值留在同一行,对 RAG 更友好(_parse_excel_to_markdown:381) |
.xml | 专用 XMLChunker | 解析和切块合成一步,保留层级(core/parser/xml_chunker.py:26) |
| 视频 | VideoParser 抽帧 + 可选转写 | 帧描述与字幕带时间戳(_parse_video:430) |
| 其他复杂格式 | Docling(本地或远程 GPU) | 通用兜底(_parse_document_local:527) |
切块器:自己实现的递归字符切分
项目没有依赖 LangChain,而是自带了一份 RecursiveCharacterTextSplitter(core/parser/morphik_parser.py:54),分隔符优先级是 ["\n\n", "\n", ". ", " ", ""]——先按段落切,切不动再按行、按句、按词、最后按字符。默认块大小 6000 字符、重叠 300(见 morphik.toml 的 [parser])。
还有一个可选的 ContextualChunker(core/parser/morphik_parser.py:110),它对每个块额外调一次 LLM,让模型写一句"这块在全文里处于什么位置",拼在块内容前面再去做嵌入:
context = self._situate_context(text, chunk.content)
content = f"{context}; {chunk.content}"
这是 Anthropic "contextual retrieval" 那套做法的实现。默认关闭(use_contextual_chunking = false),因为每块一次 LLM 调用很贵。
4. 图片路:一页一张图
入口是 _create_chunks_multivector(core/services/ingestion_service.py:1444)。它按 MIME 分派:
MIME 判断
|
+-- image/* --> 缩到宽 256、转 JPEG 质量 70,当成一个 chunk
+-- application/pdf --> _process_pdf_for_colpali,一页一个 chunk
+-- Word / PPT --> 先用 LibreOffice 转 PDF,再走 PDF 路
+-- Excel --> 不转图!用文本块走 ColPali 的文本通道
+-- 其他 --> 退回文本块,标 is_image = False
Excel 那条分支的注释很直白(core/services/ingestion_service.py:1549-1551):表格数据用文本检索效果好得多,没必要转图。
PDF 渲染:密度自适应
_process_pdf_for_colpali(core/services/ingestion_service.py:1562)会先估一个"每页多少 MB":
density_mb_per_page = file_size_mb / page_count
HIGH_DENSITY_THRESHOLD_MB = 1.0
HIGH_DENSITY_BATCH_SIZE = 2
超过 1 MB/页说明页面里塞满大图或矢量图,一次性全渲染会 OOM,于是改用 _render_pdf_with_pymupdf_batched 每次只渲 2 页。低密度文档才一次渲完。渲染 DPI 由 [pdf] colpali_pdf_dpi 控制,默认 150。
PyMuPDF 失败还有 pdf2image 兜底(第 1605-1628 行),两个都失败才抛 PdfConversionError。
白页占位:为了页号不错位
单页渲染失败时,代码不是跳过这一页,而是塞一张白图(core/services/ingestion_service.py:1401):
@staticmethod
def _placeholder_page_png() -> bytes:
"""White stand-in for a page that fails to render.
Every page must produce exactly one chunk so chunk numbers stay aligned
with page numbers; dropping a page would shift all later chunks.
"""
这是"页即块"这个约定的守护措施:引用要能说出"第 37 页",chunk_number 就必须和页号严格一一对应。丢一页,后面所有页的编号全错。
5. 降级阶梯:三层补救
如果走到某一步一个块都没产出,worker 不会直接失败,而是逐级降级。这是整个摄取里工程味最重的一段(core/workers/ingestion_worker.py:880-1018)。
图片路产出 0 块?
|
| [降级一] 回头跑一次普通文本解析 parse_file_to_text
| 拿到文本就切块,标 is_image=False 继续走
v
还是 0 块?
|
| [降级二] 跑 parse_file_to_text_deep:
| Office 文件先用 LibreOffice 转 PDF,
| 再用打开 EasyOCR + accurate 表格模式的 Docling 重解析
v
仍然 0 块?
|
| [降级三] 不失败。标记 content_extraction_status =
| "no_content_extracted",文档保留,只是搜不到
v
status = completed(带警告)
第三级的措辞值得抄(core/workers/ingestion_worker.py:995-998):
文档已成功保存,但在能提取出内容之前将不可被检索。
"深度解析"具体贵在哪,看 _build_deep_pdf_converter(core/parser/morphik_parser.py:346):打开 OCR、打开表格结构识别并调成 accurate 模式、生成页面与图片位图、images_scale = 2.0。所以它只在正常路径全军覆没时才跑。
6. 存储阶段:按批流式处理
ColPali 路径下,"嵌入全部页 → 再统一入库"会让几百页的 float 数组同时驻留内存。代码改成了批内闭环(core/workers/ingestion_worker.py:1161-1213):
for 每 16 页:
嵌入这 16 页 (embed_for_ingestion)
生成 chunk 对象 (start_index = 全局偏移,保证页号连续)
立刻写向量库 (store_embeddings)
释放
批大小来自 [worker] colpali_store_batch_size,默认 16。注释写得很清楚:"Store this batch immediately to release memory pressure"。
_create_chunk_objects 接受 start_index 参数(core/services/ingestion_service.py:1207),这就是分批之后页号仍然连续的原因。
重试前先清场
arq 配置了 max_tries = 5、retry_jobs = True(core/workers/ingestion_worker.py:2064-2069)。重试意味着整个任务从头再跑一遍——如果上次已经写了一半的 chunk 进去,插入式的存储(Postgres 多向量、pgvector)就会留下重复行。
处理办法在第 1090-1091 行:
is_retry_attempt = int(ctx.get("job_try") or 1) > 1
if doc.chunk_ids or is_retry_attempt:
# 按 document_id 删掉所有已有 chunk 再重写
注释点出了关键:删除按 document_id 而不是按 chunk_ids,所以即使上次失败得连 chunk_ids 都没来得及持久化,也照样清得干净。
7. 关键细节与坑
- 解析用的文件名不是用户给的文件名。 第 663 行
parse_filename = os.path.basename(file_key),用的是存储键推出来的名字。注释解释:UI 传上来的original_filename常常是.pdf,但实际存的可能是预抽取的.pdf.txt,用错扩展名会误导解析器。 - worker 每个任务新建一个 IngestionService。 第 579 行。目的是避免并发任务之间串号(不同 app 的向量库实例被互相覆盖)。但底层的向量库实例是按配置缓存共享的(
_get_worker_colpali_store:116,带双检锁)。 - 页数是估出来的,不是数出来的。 计费口径见
core/limits_utils.py:57的estimate_pages_by_chars:4 字符 = 1 token,630 token = 1 页。但 ColPali 路径下会用真实图片块数覆盖这个估值(core/workers/ingestion_worker.py:1022-1024)。 - 文档更新走的是另一条队列路径,见
queue_document_update(core/services/ingestion_service.py:704),不要和首次摄取混淆。
8. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 摄取 HTTP 入口 | core/routes/ingest.py | ingest_file、ingest_text、batch_ingest_files |
| 落库 + 上传 + 入队 | core/services/ingestion_service.py | ingest_file_content、_build_ingestion_job_payload |
| 失败标记 | core/services/ingestion_service.py | _mark_document_failed |
| worker 主流程 | core/workers/ingestion_worker.py | process_ingestion_job |
| 进度上报 | core/workers/ingestion_worker.py | update_document_progress |
| 向量库实例缓存 | core/workers/ingestion_worker.py | _get_worker_colpali_store |
| arq 运行参数 | core/workers/ingestion_worker.py | WorkerSettings |
| 解析总入口 | core/parser/morphik_parser.py | MorphikParser.parse_file_to_text |
| 深度解析兜底 | core/parser/morphik_parser.py | parse_file_to_text_deep、_build_deep_pdf_converter |
| 自实现切块器 | core/parser/morphik_parser.py | RecursiveCharacterTextSplitter |
| 上下文增强切块 | core/parser/morphik_parser.py | ContextualChunker |
| XML 结构化切块 | core/parser/xml_chunker.py | XMLChunker |
| 页面图片块生成 | core/services/ingestion_service.py | _create_chunks_multivector、_process_pdf_for_colpali |
| 白页占位 | core/services/ingestion_service.py | _placeholder_page_png |
| ColPali 原生格式判定 | core/storage/utils_file_extensions.py | is_colpali_native_format |
| 页数估算 | core/limits_utils.py | estimate_pages_by_chars |