非结构化文档 ETL 与多模态数据

企业 80% 的数据躺在 PDF、扫描件、Office 文档、邮件与音视频里,它们无法被 SQL 直接查询。本文讲清非结构化 ETL 的完整链路:文档解析工具选型与 PDF 三种形态的差异、OCR 与版面分析、表格抽取、分块策略与重叠控制、元数据与血缘注入、嵌入生成与向量入库、增量重处理与幂等设计,以及质量评估、成本控制与常见踩坑清单。

引言

数仓与湖仓擅长处理结构化数据:订单、日志、指标,有 Schema、可 SQL 查询。但企业里真正体量最大的是非结构化内容——合同、发票、研报、病历、工单、邮件、会议录音。它们躺在文件服务器与对象存储里,既不可查询,也不可关联,更谈不上进入模型。

非结构化文档 ETL(Unstructured Document ETL)就是把这类内容转化为可检索、可关联、可治理的数据资产:解析出文本与版面结构,抽取关键字段,切分成语义单元,生成嵌入并入库,同时保留完整的元数据与血缘。本文按这条链路逐步拆解工程实现。

一、非结构化数据的范围与挑战

1.1 数据形态

类别格式解析难点
数字原生 PDFPDF(含文字层)版面还原、阅读顺序、表格结构
扫描件 PDF图片型 PDF必须 OCR,质量依赖图像
Office 文档docx / xlsx / pptx样式与层级、嵌套表格、批注
邮件eml / msg线程关系、附件递归
网页HTML正文抽取、去导航与广告
图像jpg / png需视觉模型或 OCR
音频 / 视频mp3 / mp4需 ASR 转写,时间轴对齐

1.2 与传统 ETL 的根本差异

结构化 ETL:Schema 已知 → 映射 → 转换 → 加载
非结构化 ETL:Schema 未知 → 推断结构(解析/抽取)→ 归一化 → 分块 → 嵌入 → 加载

三个额外成本:

  1. 解析是概率性的:同一份 PDF,不同解析器给出的文本顺序可能不同,不存在"绝对正确"。
  2. 结果是多模态的:除文本外还有版面坐标、表格、图片、音频时间轴。
  3. 质量难以自动断言:结构化数据的 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

三条工程纪律:

  1. 嵌入要幂等:以 chunk 内容哈希为键,内容未变不重复调用嵌入 API。
  2. 模型版本要落库:换嵌入模型必须全量重嵌入,向量不可跨模型混用。
  3. 失败要可续传:批量作业按 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,再选带版面能力的解析器保住结构,最后用合理的分块与嵌入策略把内容变成可检索的资产。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 湖仓访问控制与权限治理
  2. 流式 SQL:Flink SQL 与 ksqlDB
  3. Flink 状态后端与 Checkpoint 调优