跳到主要内容

数据截至 (上游 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:552ingest_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:448process_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/*
PDFapplication/pdf
医学影像application/dicom
Word...wordprocessingml.documentapplication/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/.xlsmopenpyxl 转 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 = 5retry_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:57estimate_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.pyingest_fileingest_textbatch_ingest_files
落库 + 上传 + 入队core/services/ingestion_service.pyingest_file_content_build_ingestion_job_payload
失败标记core/services/ingestion_service.py_mark_document_failed
worker 主流程core/workers/ingestion_worker.pyprocess_ingestion_job
进度上报core/workers/ingestion_worker.pyupdate_document_progress
向量库实例缓存core/workers/ingestion_worker.py_get_worker_colpali_store
arq 运行参数core/workers/ingestion_worker.pyWorkerSettings
解析总入口core/parser/morphik_parser.pyMorphikParser.parse_file_to_text
深度解析兜底core/parser/morphik_parser.pyparse_file_to_text_deep_build_deep_pdf_converter
自实现切块器core/parser/morphik_parser.pyRecursiveCharacterTextSplitter
上下文增强切块core/parser/morphik_parser.pyContextualChunker
XML 结构化切块core/parser/xml_chunker.pyXMLChunker
页面图片块生成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.pyis_colpali_native_format
页数估算core/limits_utils.pyestimate_pages_by_chars