引言
数仓与湖仓擅长处理结构化数据:订单、日志、指标,有 Schema、可 SQL 查询。但企业里真正体量最大的是非结构化内容——合同、发票、研报、病历、工单、邮件、会议录音。它们躺在文件服务器与对象存储里,既不可查询,也不可关联,更谈不上进入模型。
非结构化文档 ETL(Unstructured Document ETL)就是把这类内容转化为可检索、可关联、可治理的数据资产:解析出文本与版面结构,抽取关键字段,切分成语义单元,生成嵌入并入库,同时保留完整的元数据与血缘。本文按这条链路逐步拆解工程实现。
一、非结构化数据的范围与挑战
1.1 数据形态
| 类别 | 格式 | 解析难点 |
|---|---|---|
| 数字原生 PDF | PDF(含文字层) | 版面还原、阅读顺序、表格结构 |
| 扫描件 PDF | 图片型 PDF | 必须 OCR,质量依赖图像 |
| Office 文档 | docx / xlsx / pptx | 样式与层级、嵌套表格、批注 |
| 邮件 | eml / msg | 线程关系、附件递归 |
| 网页 | HTML | 正文抽取、去导航与广告 |
| 图像 | jpg / png | 需视觉模型或 OCR |
| 音频 / 视频 | mp3 / mp4 | 需 ASR 转写,时间轴对齐 |
1.2 与传统 ETL 的根本差异
结构化 ETL:Schema 已知 → 映射 → 转换 → 加载
非结构化 ETL:Schema 未知 → 推断结构(解析/抽取)→ 归一化 → 分块 → 嵌入 → 加载
三个额外成本:
- 解析是概率性的:同一份 PDF,不同解析器给出的文本顺序可能不同,不存在"绝对正确"。
- 结果是多模态的:除文本外还有版面坐标、表格、图片、音频时间轴。
- 质量难以自动断言:结构化数据的
not_null在这里没有对应物,需要构造专门的评估集。
二、文档解析
2.1 解析工具选型
| 工具 | 擅长 | 局限 |
|---|---|---|
| PyMuPDF (fitz) | 快速提取文字与坐标 | 表格结构弱 |
| pdfplumber | 精确坐标、表格线框 | 慢 |
| Apache Tika | 格式覆盖最广(千余种) | 只出纯文本,丢版面 |
| Unstructured | 元素级切分(Title/NarrativeText/Table) | 依赖较重 |
| Docling | 版面模型 + 表格结构(TableFormer) | 计算开销大 |
| 商业 API | 精度最高,含表格与手写 | 成本、数据出境 |
选型原则:先判断文档是数字原生还是扫描件——这决定要不要 OCR,也决定解析器的选择。
import fitz # PyMuPDF
doc = fitz.open("contract.pdf")
page = doc[0]
text = page.get_text("text")
# 关键判断:文字层是否存在(决定是否走 OCR)
if len(text.strip()) < 50 and len(page.get_images()) > 0:
route = "ocr"
else:
route = "text_layer"
print(f"pages={len(doc)} route={route} chars={len(text)}")
2.2 PDF 的三种形态
1. 数字原生(Text-based) 有文字层,直接提取,顺序可能乱
2. 扫描件(Image-based) 无文字层,必须 OCR
3. 混合型 部分页有文字层,部分没有
混合型最容易被忽略:一个 200 页的文档里夹了 5 页扫描件,如果只按"首页有无文字层"判断,这 5 页会静默丢失。正确做法是逐页判断:
for i, page in enumerate(doc):
t = page.get_text("text").strip()
if len(t) < 20:
pages_needing_ocr.append(i)
2.3 版面分析与阅读顺序
多栏排版、页眉页脚、脚注会让简单的文本提取产生错乱顺序。用坐标聚类恢复阅读顺序:
import pdfplumber
with pdfplumber.open("report.pdf") as pdf:
page = pdf.pages[0]
# 按 y 坐标分行,再按 x 坐标排序,近似还原阅读顺序
words = page.extract_words(use_text_flow=False, keep_blank_chars=False)
lines = {}
for w in words:
key = round(w["top"] / 5) # 5pt 容差归并同一行
lines.setdefault(key, []).append(w)
ordered = [
" ".join(x["text"] for x in sorted(lines[k], key=lambda w: w["x0"]))
for k in sorted(lines)
]
对表格密集的文档,优先用带版面模型的解析器(Docling 的 TableFormer、Unstructured 的 hi_res 策略),它们会把表格输出为结构化行,而不是把单元格文本混进正文。
2.4 OCR
# Tesseract:轻量、离线、中文需装语言包
tesseract scan.png out -l chi_sim+eng --psm 6
# PaddleOCR:中文场景精度更高,支持版面分析
python -m paddleocr --image_dir ./pages --lang ch --use_angle_cls true
OCR 的质量决定下游一切。工程上的关键措施:
| 措施 | 说明 |
|---|---|
| 预处理 | 二值化、去噪、倾斜校正(deskew)、300 DPI 以上 |
| 语言包 | 中英混排必须同时加载 chi_sim+eng |
| 置信度阈值 | 低于阈值(如 0.6)的片段标记待人工复核 |
| 版面对齐 | 用 OCR 结果的坐标还原表格与段落 |
2.5 解析层要可替换
解析器迭代速度极快,一年内主力工具可能换两轮。因此解析层必须做成可插拔的适配器,统一输出同一种中间表示:
from dataclasses import dataclass, field
from typing import Protocol
@dataclass
class Block:
type: str # title / narrative_text / table / figure
text: str
page: int
bbox: tuple = field(default_factory=tuple)
level: int = 0 # 标题层级
class Parser(Protocol):
name: str
version: str
def parse(self, path: str) -> list[Block]: ...
# 不同解析器实现同一接口,产物统一为 list[Block]
PARSERS = {"pymupdf": PyMuPDFParser(), "docling": DoclingParser(), "tika": TikaParser()}
这样切换解析器只需改配置,并把 parser + version 写进元数据,出问题时能精确定位是哪一批解析产物受影响、需要重跑。
三、分块策略
分块(Chunking)决定了检索的粒度,是文档 ETL 里最影响最终效果的环节。
3.1 三种分块方式
| 方式 | 做法 | 适用 |
|---|---|---|
| 固定长度 | 按 token 数硬切,带重叠 | 通用兜底 |
| 结构感知 | 按标题层级、段落、表格边界切 | 有明确结构的文档 |
| 语义分块 | 按句子嵌入相似度断点 | 长叙述型文本 |
from langchain_text_splitters import RecursiveCharacterTextSplitter
splitter = RecursiveCharacterTextSplitter(
chunk_size=800, # 目标 token 数(约 600~1000 字)
chunk_overlap=120, # 15% 重叠,避免语义在边界被截断
separators=["\n## ", "\n### ", "\n\n", "\n", "。"], # 优先级从高到低
length_function=len,
)
chunks = splitter.split_text(markdown_text)
3.2 参数与权衡
| 参数 | 偏小 | 偏大 |
|---|---|---|
| chunk_size | 检索精准但上下文不足 | 上下文完整但噪声多、嵌入被稀释 |
| chunk_overlap | 边界信息丢失 | 冗余、存储与嵌入成本上升 |
| 分隔符优先级 | 破坏结构 | 可能产出超长块 |
经验起点:中文 600~1000 字、重叠 10%~15%,且每个 chunk 必须携带标题路径(如 合同 > 第 3 条 > 付款方式),否则检索到片段却不知道它属于哪一节。
def chunk_with_breadcrumb(doc_sections):
out = []
for sec in doc_sections:
breadcrumb = " > ".join(sec["path"]) # 标题路径
for i, c in enumerate(splitter.split_text(sec["text"])):
out.append({
"text": f"[{breadcrumb}]\n{c}", # 拼进正文,提升召回
"metadata": {"breadcrumb": breadcrumb, "chunk_index": i,
"doc_id": sec["doc_id"], "page": sec["page"]},
})
return out
四、元数据与血缘
非结构化数据进入数据平台后,没有元数据就等于不可用。每个 chunk 至少携带:
{
"doc_id": "contract-2026-000123",
"source_uri": "s3://docs/contracts/2026/000123.pdf",
"source_hash": "sha256:9f2c...",
"page": 7,
"breadcrumb": "合同 > 第 3 条 > 付款方式",
"chunk_index": 12,
"parser": "docling@2.1.0",
"ocr_used": false,
"extracted_at": "2026-10-07T19:18:00+08:00",
"acl": ["group:legal", "group:finance"],
"language": "zh"
}
source_hash 是增量处理的基础:文档内容未变则跳过解析与嵌入,直接复用已有向量。acl 字段让检索时能做权限过滤,避免越权召回——这一点在合规场景里是硬要求,权限与脱敏模型可参考 https://plumephp.com/data-security-privacy-compliance/。
血缘层面,文档解析应作为一个节点接入统一血缘,向上游追到对象存储路径,向下游追到向量集合与字段抽取结果,与结构化链路的血缘体系一致(见 https://plumephp.com/data-catalog-lineage/)。
五、嵌入生成与入库
5.1 批量嵌入与限流
import hashlib
from tenacity import retry, wait_exponential, stop_after_attempt
@retry(wait=wait_exponential(multiplier=1, min=1, max=30), stop=stop_after_attempt(5))
def embed_batch(texts, model="bge-m3", batch_size=64):
vectors = []
for i in range(0, len(texts), batch_size):
vectors.extend(client.embed(model, texts[i:i + batch_size]))
return vectors
def needs_reembed(chunk, existing_hash):
h = hashlib.sha256(chunk["text"].encode()).hexdigest()
return h != existing_hash
三条工程纪律:
- 嵌入要幂等:以 chunk 内容哈希为键,内容未变不重复调用嵌入 API。
- 模型版本要落库:换嵌入模型必须全量重嵌入,向量不可跨模型混用。
- 失败要可续传:批量作业按 doc_id 记录进度,中断后从未完成处继续。
5.2 入库
向量库与检索索引的选型、索引参数(HNSW 的 M / ef_construction)与混合检索的落地,详见 https://plumephp.com/data-vector-database-rag-pipeline/。此处只强调一点:向量集合要按"模型版本 + 分块策略版本"分区,而不是把不同版本的向量混在一个集合里。
-- 元数据表(PostgreSQL):支撑过滤与治理
CREATE TABLE doc_chunks (
chunk_id TEXT PRIMARY KEY,
doc_id TEXT NOT NULL,
content TEXT NOT NULL,
content_hash TEXT NOT NULL,
breadcrumb TEXT,
page INT,
embed_model TEXT NOT NULL,
chunk_version TEXT NOT NULL,
acl TEXT[],
created_at TIMESTAMPTZ DEFAULT now()
);
CREATE INDEX ON doc_chunks (doc_id);
CREATE INDEX ON doc_chunks USING gin (acl);
六、管道编排与增量处理
6.1 分层管道
Bronze:原始文件落地(对象存储 + 内容哈希清单)
Silver:解析产物(文本 + 版面 JSON + 表格 CSV)
Gold:分块与嵌入(chunk 表 + 向量集合)
分层的意义在于可重放:解析器升级时只需重跑 Silver → Gold,无需重新拉取原始文件;分块策略调整时只需重跑 Gold。这与结构化 ETL 的 Bronze/Silver/Gold 分层完全同构,设计原则见 https://plumephp.com/etl-elt-design/。
6.2 增量与幂等
def sync_documents(source_root, manifest_table):
for uri, stat in scan(source_root):
h = sha256_of(uri)
prev = manifest_table.get(uri)
if prev and prev["hash"] == h:
continue # 内容未变,跳过
yield {"uri": uri, "hash": h, "action": "parse"}
编排层用"文档哈希 + 解析器版本 + 分块版本"三元组做幂等键:任一变化才触发重处理。批量解析作业通常按文档粒度并行,单个文档失败不阻塞整批。
七、质量与成本
7.1 质量评估
非结构化 ETL 的质量无法靠 not_null 保证,需要构造评估集:
| 指标 | 定义 | 采集方式 |
|---|---|---|
| 解析成功率 | 成功产出文本的文档占比 | 管道埋点 |
| 字符损失率 | 解析字符数 / 原文字符数 | 抽样比对 |
| OCR 置信度均值 | 低于阈值的页占比 | OCR 输出统计 |
| 抽取字段准确率 | 人工标注集上的 F1 | 定期人工评估 |
| 检索命中率 | 评估问题集中 top-k 命中率 | 评估集回放 |
7.2 成本控制
成本大头 = OCR 计算 + 嵌入 API 调用 + 向量存储 + 重处理
优化顺序:
1. 哈希去重,跳过未变文档(收益最大)
2. 先判断有无文字层,能不提 OCR 就不提
3. 嵌入按内容哈希缓存,跨批次复用
4. 向量量化(PQ/SQ)降低存储,牺牲少量召回
5. 分块粒度合理,避免过细导致嵌入次数翻倍
八、踩坑清单
| 坑 | 表现 | 修法 |
|---|---|---|
| 只按首页判断是否 OCR | 混合型文档静默丢页 | 逐页判断文字层 |
| 分块不带标题路径 | 检索到片段不知出处 | chunk 前置 breadcrumb |
| 向量跨模型混用 | 检索结果诡异 | 集合按模型版本隔离 |
| 无内容哈希 | 每次全量重嵌入,成本失控 | 以哈希做幂等键 |
| 忽略阅读顺序 | 多栏文档文本错乱 | 坐标聚类恢复顺序 |
| 表格被压成纯文本 | 关键数值丢失结构 | 用版面模型输出结构化表格 |
| 元数据缺 ACL | 越权召回敏感文档 | 入库即写权限标签 |
| 无评估集 | 效果退化无从察觉 | 固定评估集 + 定期回放 |
小结
非结构化文档 ETL 的难点不在"能不能解析",而在可重放、可治理、可评估。可重放靠 Bronze/Silver/Gold 分层与内容哈希;可治理靠完整元数据(来源、标题路径、权限、解析器版本)与血缘接入;可评估靠固定评估集与解析质量指标。技术上,先判断文档形态决定是否 OCR,再选带版面能力的解析器保住结构,最后用合理的分块与嵌入策略把内容变成可检索的资产。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。