数据截至 (上游 commit c4b5ed6202d6)
流水线层:声明式直出 vs 多线程分阶段
30 秒导读: 后端(02)把字节变成"可处理的形态"后,流水线(pipeline) 负责把它真正加工成
DoclingDocument。所有流水线都走同一条骨架——搭建 → 装配 → 富化 → 定状态; 真正拉开差距的是"搭建"这一步:能直出的格式一行调用就完事,PDF 则要拆成 OCR/版面/表格等多个阶段、 用多线程流水并行跑。本章讲这条骨架和三种主流水线(Simple / StandardPdf / Vlm)的取舍。
本章位置:后端把原始输入变成 backend 对象后,DocumentConverter 按格式挑一条流水线来跑
(挑选逻辑见 01 入口与分派)。各阶段内部的模型(版面、表格、OCR 具体怎么算)
留给 04;DoclingDocument 怎么拼、怎么导出给 RAG 留给 05。
1. 这是什么(零基础也能懂)
一句话定义: 流水线是"把一份输入文档加工成统一 DoclingDocument 的装配车间"——
它规定了先做什么、后做什么、出错怎么办。
为什么需要它。 不同格式的加工难度天差地别:
- 一个 Markdown / HTML 文件,后端本身就能直接吐出结构化文档——几乎不用"加工"。
- 一个扫描版 PDF,得先渲染每一页、跑 OCR 认字、跑版面模型分块、跑表格模型还原表格、 再按阅读顺序拼起来——十几个步骤,还要顶得住上百页的量。
如果每种格式各写一套流程,代码会散架。Docling 的做法是:用一条统一骨架框住所有流水线, 把"每种格式的差异"收敛到骨架里的少数几个可重写方法。于是:
- 简单格式 → 重写"搭建"这一步,一行直出。
- PDF → 重写"搭建"这一步,塞进一整套多线程分阶段流水。
一句话直觉: 把流水线想成工厂流水线的总控。传送带的走向(骨架)是固定的; 不同产品(格式)只是换掉传送带上的加工机器。
用起来什么样。 使用者几乎感觉不到流水线的存在——DocumentConverter 会替你选好:
# 示意,非源码
from docling.document_converter import DocumentConverter
conv = DocumentConverter()
res = conv.convert("scan.pdf") # 背后跑 StandardPdfPipeline(多线程分阶段)
res = conv.convert("notes.md") # 背后跑 SimplePipeline(后端直出)
print(res.document.export_to_markdown())
重点看:同一个 API,背后是两条完全不同的流水线;这一层的价值就是"把差异藏起来"。
2. 顶层全景(它大概怎么转)
2.1 统一骨架:每条流水线都跑这四步
不管哪种流水线,对外只有一个入口 execute(),它固定按四步走。真源码在
docling/pipeline/base_pipeline.py:65(BasePipeline.execute):
BasePipeline.execute(in_doc)
│
┌───────────────────┼────────────────────┐
│ ① _build_document 搭建:原始 → 页/文档骨架 │
│ ② _assemble_document 装配:页拼成 DoclingDocument │
│ ③ _enrich_document 富化:对文档元素做二次加工 │
│ ④ _determine_status 定状态:SUCCESS / PARTIAL… │
└───────────────────┬────────────────────┘
│ (无论成败)
finally: _unload 卸载后端、释放资源
怎么读这张图: 从上到下是严格顺序;①②只管"搭出文档结构",从③起只依赖 conv_res.document
(源码在 base_pipeline.py:77 明确注释了这个边界)。整段包在 try/except/finally 里——
异常→状态判 FAILURE,finally 里 必定 _unload 释放资源。
四步各自的职责:
| 步骤 | 方法 | 干什么 | 谁必须实现 |
|---|---|---|---|
| ① 搭建 | _build_document | 把后端产物变成页列表或直接的文档 | @abstractmethod,每条流水线必写 |
| ② 装配 | _assemble_document | 把逐页结果拼成一个 DoclingDocument | 默认空实现,PDF/VLM 重写 |
| ③ 富化 | _enrich_document | 遍历文档元素跑 enrichment_pipe | 基类通用实现,一般不重写 |
| ④ 定状态 | _determine_status | 决定 SUCCESS / PARTIAL / FAILURE | @abstractmethod,每条流水线必写 |
2.2 一个关键收尾:PARTIAL_SUCCESS 与"有错不算成功"
_determine_status 定完状态后,execute 还有一道兜底(base_pipeline.py:82):
只要 conv_res.errors 非空,就绝不报 SUCCESS,自动降级为 PARTIAL_SUCCESS。
# 摘自 base_pipeline.py:83-87(execute 内)
conv_res.status = self._determine_status(conv_res)
if conv_res.status == ConversionStatus.SUCCESS and conv_res.errors:
conv_res.status = ConversionStatus.PARTIAL_SUCCESS
这条规则贯穿全书:某一 页 OCR 挂了、某页 VLM 输出被截断,文档整体还能出——但状态诚实地标成
"部分成功",错误明细进 conv_res.errors。
2.3 三种主流水线一览
| 流水线 | 文件 | 适用输入 | 搭建方式 |
|---|---|---|---|
SimplePipeline | simple_pipeline.py | MD/HTML/DOCX 等声明式后端 | 后端一步 convert() 直出 |
StandardPdfPipeline | standard_pdf_pipeline.py | PDF / 图片 | 多线程分阶段流水(本章重点) |
VlmPipeline | vlm_pipeline.py | PDF(视觉大模型路线) | 逐页喂给 VLM,整页转 DocTags/MD |
AsrPipeline | asr_pipeline.py | 音频/视频 | Whisper 转写成带时间戳的文本 |
3. 类层次:骨架如何一层层被特化
理解流水线的关键,是看清继承链上每一层新增了什么。
BasePipeline (ABC) ← execute 骨架 + enrichment 遍历
│ base_pipeline.py:50
├── ConvertPipeline ← 装配"通用富化模型"(图片分类/描述/图表)
│ base_pipeline.py:153
│ ├── SimplePipeline ← _build_document = backend.convert()
│ │ simple_pipeline.py:16
│ ├── StandardPdfPipeline ← 自带多线程 _build_document(本章重点)
│ │ standard_pdf_pipeline.py:567
│ └── PaginatedPipeline ← 按 page_batch 逐页跑 build_pipe
│ base_pipeline.py:231
│ └── VlmPipeline ← build_pipe = 整页 VLM
│ vlm_pipeline.py:83
└── AsrPipeline ← 直接继承 BasePipeline(不需要富化/分页)
asr_pipeline.py:30
注意:
StandardPdfPipeline直接继承ConvertPipeline而不走PaginatedPipeline—— 它用自己的线程化_build_document取代了逐页循环。PaginatedPipeline现在主要服务VlmPipeline(和已弃用的 legacy PDF 流水)。
3.1 BasePipeline——骨架 + 富化遍历
除了 §2 的 execute 骨架,基类还实现了通用的 _enrich_document
(base_pipeline.py:107)。它做的事很固定:
- 遍历
conv_res.document的所有元素(iterate_items); - 对
enrichment_pipe里的每个模型,用model.prepare_element挑出它关心的元素; - 按
elements_batch_size分批喂给模型。
关键一行是那句"必须耗尽"的注释(base_pipeline.py:124):模型返回的是生成器,
for element in model(...) 必须走完,否则副作用不落地。
3.2 ConvertPipeline——预置"通用富化模型"
ConvertPipeline(base_pipeline.py:149)在构造时就把三类跨后端通用的富化模型装进
enrichment_pipe:
| 富化模型 | 干什么 | 源码符号 |
|---|---|---|
| 图片分类 | 判断插图是照片/图标/图表… | DocumentPictureClassifier |
| 图片描述 | 给插图生成文字描述(可接远程模型) | _get_picture_description_model |
| 图表抽取 | 把图表还原成结构化数据 | ChartExtractionModelGraniteVision / …V4 |
一个细节:图表抽取依赖图片分类,所以构造函数里用局部变量把
do_picture_classification 或上 do_chart_extraction(base_pipeline.py:156),
避免直接改动共享的 pipeline_options(改了会影响其哈希、破坏缓存)。
3.3 PaginatedPipeline——逐页跑 build_pipe
PaginatedPipeline(base_pipeline.py:238)是"顺序分页"的流水线基类。它的
_build_document(base_pipeline.py:251)骨架是:
按 page_batch_size 把页切成一批批(默认 4,settings.py:32)
每批:
1. initialize_page 给每页挂上 page backend、量尺寸
2. _apply_on_pages 让 build_pipe 里每个 model 依次处理这一批页
3. 逐页清理 清 image_cache、unload page backend
每批结束检查 document_timeout;超时→PARTIAL_SUCCESS,break
两个要点:
_apply_on_pages(base_pipeline.py:243)就是 build_pipe 的执行器——把page_batch依次穿过每个 build 模型:page_batch = model(conv_res, page_batch)。- 超时是"批粒度"的:累计耗时超过
document_timeout就停,已处理的页保留, 未处理的页最后被过滤掉(base_pipeline.py:340剔除size is None的页)。
_determine_status(base_pipeline.py:360)则扫描每一页的 backend:任何一页 backend 失效,
就追加一条 BACKEND_FAILURE 错误并降级 PARTIAL_SUCCESS。
4. build_pipe vs enrichment_pipe(一定要分清)
这是本层最容易混的两个"管道"。它们跑在不同的对象、不同的阶段上。
build_pipe | enrichment_pipe | |
|---|---|---|
| 作用对象 | Page 对象(逐页) | DoclingDocument 的元素(NodeItem) |
| 运行阶段 | ① 搭建期(_build_document) | ③ 富化期(_enrich_document) |
| 目的 | 把一页"看懂":OCR/版面/表格/整页 VLM | 对成品文档做二次加工:图片描述、图表抽取、代码/公式识别 |
| 谁填充它 | 各子类构造函数 | ConvertPipeline 预置 + 子类追加 |
一句话:build_pipe 负责"从像素到结构";enrichment_pipe 负责"结构好了再补料"。
时序上,富化永远在装配之后——因为富化操作的是已经拼好的 conv_res.document,
而不是零散的页(见 §2.1 的边界注释)。
5. SimplePipeline——声明式后端直出
这是最简单的一条,适合本身就能吐结构化文档的格式(Markdown、HTML、DOCX…)。
它对"搭建"的实现几乎没有内容——直接调用后端的 convert()
(simple_pipeline.py:26,SimplePipeline._build_document):
# 摘自 simple_pipeline.py:39-40
with TimeRecorder(conv_res, "doc_build", scope=ProfilingScope.DOCUMENT):
conv_res.document = conv_res.input._backend.convert()
前提是后端必 须是 DeclarativeDocumentBackend(能直出文档的后端,见 02);
否则直接抛 RuntimeError。_determine_status 也简单——没有别的可评估,直接 SUCCESS
(simple_pipeline.py:43)。
为什么它没有 build_pipe。 因为"看懂一页"的活儿全在后端里做完了,流水线无须再逐页加工——
它只借用基类的 execute 骨架和 ConvertPipeline 的富化能力(仍会跑图片分类/描述等)。
6. StandardPdfPipeline——多线程分阶段(本章重点)
PDF 是最重的路径,也是这个文件里工程密度最高的地方。它没有用 PaginatedPipeline 的
顺序分页,而是自己实现了一套多线程流水(standard_pdf_pipeline.py:578)。
6.1 为什么要线程化
一份 PDF 的每页要顺序经过 OCR → 版面 → 表格 → 装配。如果一页一页串行走,慢阶段会拖死整体。 更好的做法是让各阶段像工厂流水线一样并行:第 1 页在跑表格时,第 3 页可以同时在跑 OCR。
Docling 的方案:每个阶段一个专属线程,阶段之间用有界队列连起来,页像零件一样在传送带上流动。
producer 线程 每个阶段 = 1 个线程,阶段间 = 有界队列
逐页 load_page (队列满则上游阻塞 = 背压)
│
▼
[preprocess] ─q─▶ [ocr] ─q─▶ [layout] ─q─▶ [table] ─q─▶ [assemble] ─q─▶ output_q
│
主线程 drain 汇总
怎么读:横向是数据流向,
q是ThreadedQueue。每个方框独立跑在自己的线程里, 谁先干完谁就把结果丢进下游队列,不必等别人。
6.2 三块积木
积木 A:ThreadedQueue——有界队列 + 背压 + close 传播(standard_pdf_pipeline.py:155)