038AI Agent阅读记录
第 038 卷

从零开始 AI Agent 实战(九):文档导入,解析、清洗与重叠窗口切分

从 PDF、DOCX、Markdown、TXT 提取带标题层级的文本,采用重叠窗口切分并异步写入 pgvector,提供 202 与状态轮询接口。

乌漆嘛黑和 Ahri
第 038 期

文档导入:解析、清洗与重叠窗口切分

先让 80 页 PDF 把接口拖死

上一篇把群聊导成了协议包,但它还不是语料。这一篇先放下群聊,接另一个来源:上传的 PDF、Word 和 TXT。

顺序是有意的。官方手册有完整的标题层级,是两个来源里规整得多的那个,适合先把「解析 → 切分 → 向量 → 入库」这条链路打通;群聊那半留到下一篇——到时候你会发现,这一篇的切分方式对它完全无效

第一版接口很容易写成这样:

@router.post('/documents')
async def upload(file: UploadFile, user: User = Depends(current_user)):
    raw = await file.read()
    blocks = parse(raw, file.filename)          # 解析
    chunks = chunk(blocks)                       # 切分
    vectors = await embed([c.text for c in chunks])   # 逐批调 embedding
    await chunks_repo.insert(chunks, vectors)    # 入库
    return {'document_id': ..., 'chunks': len(chunks)}

拿一份 80 页的员工手册测一下,这个请求要跑 40 秒以上。时间花在哪很清楚:embed 是网络调用,按 batch_size=64 分批,80 页大约切出 600 个 chunk,就是 10 次串行的模型请求。

超时本身不是最糟的。最糟的是它的失败方式:

  • 浏览器 60 秒断开,用户看到 504
  • 服务端的协程没有被取消,它继续跑完,把 600 个 chunk 全写进了库
  • 用户以为失败了,重新上传一次
  • 现在库里有两份索引,检索时同一段话出现两次,还都会被算进 Top-K

一次超时换来一个静默的数据污染。这比直接报错难查得多。

正确的 HTTP 语义是把两件事分开:上传成功只代表文件已经收下,不代表索引完成

POST /documents -> 202 Accepted
{
  "document_id": "d_123",
  "status": "queued",
  "status_url": "/documents/d_123/status"
}
Preview
从字节到可检索向量,以及同步与异步的边界

图中那条竖直虚线是这一篇所有设计的出发点。线的左边必须在 2 秒内返回,能做的只有三件事:把文件写进对象存储、落一条 documents 记录、把任务扔进队列。剩下的全部搬到线右边。

边界画在入队处而不是解析之后,是因为解析本身也可能很慢。90 页的扫描版 PDF 光提取文字就要十几秒,把它留在请求里,接口延迟就成了「文档页数 × 每页解析耗时」这个没有上界的函数。

解析:把结构一起带出来

最省事的解析是 pdf.get_text() 拼成一个大字符串。但这样做,第 12 篇就没法生成引用了——用户看到一句「清理缓存目录后重启」,无法验证它出自哪一节。

所以所有格式的解析器都返回同一个结构,把位置信息一起带出来:

@dataclass
class DocumentBlock:
    text: str
    heading_path: tuple[str, ...]   # ('第四章 故障排查', '4.3 启动异常')
    page: int | None = None
    source_offset: int = 0          # 在原文中的字符偏移

heading_path 是这里最有价值的字段。它有两个用途:一是引用时能显示「出自 4.3 启动异常」,二是切分时可以把标题拼进 chunk 正文——一段孤立的「删除 %APPDATA% 下的缓存目录后重启」几乎无法被检索命中,但带上标题路径之后,它同时含有「故障排查」和「启动异常」两个关键词。

Markdown 的解析器最短,能看清整体思路:

def parse_markdown(raw: str) -> list[DocumentBlock]:
    blocks, headings, offset = [], [], 0
    for paragraph in split_markdown_blocks(raw):
        if paragraph.startswith('#'):
            level = len(paragraph) - len(paragraph.lstrip('#'))
            headings = headings[:level - 1] + [paragraph[level:].strip()]
            offset += len(paragraph)
            continue
        blocks.append(DocumentBlock(paragraph.strip(), tuple(headings), source_offset=offset))
        offset += len(paragraph)
    return blocks

headings[:level - 1] + [...] 这一行在维护一个标题栈:遇到 ## 就截断到一级,保证 heading_path 始终是当前段落的完整祖先链。

其他三种格式各有各的坑:

  • PDFpymupdf。它能给出页码,但没有标题概念——需要按字号推断层级,字号明显大于正文的行视作标题。这个判断不可能完全准确,接受它。
  • DOCXpython-docx,标题反而最好认:paragraph.style.name 就是 Heading 1Heading 2。但它没有页码,因为 Word 的分页是渲染时才确定的。
  • TXT 没有任何结构,heading_path 为空。这时文件名是唯一的上下文来源,要把它塞进 chunk。

编码是 TXT 的另一个麻烦。国内的 .txt 有相当比例是 GBK,直接 decode('utf-8') 会抛异常:

def decode_text(raw: bytes) -> str:
    for encoding in ('utf-8', 'gb18030', 'utf-16'):
        try:
            return raw.decode(encoding)
        except UnicodeDecodeError:
            continue
    return raw.decode('utf-8', errors='replace')

顺序有讲究:gb18030 放在 gbk 的位置上,因为它是 GBK 的超集,能多认一批生僻字。最后一行的 errors='replace' 是兜底——把无法识别的字节变成 <?> 也比让整个索引任务失败要好。

切分:重叠窗口是为了不切断句子

固定字符数切分不高明,但它可复现,而且不需要 tokenizer。默认 chunk_size=500overlap=80

Preview
三种切分粒度与重叠窗口的作用

重叠存在的唯一理由是边界。一段排查步骤写着「……需先结束残留进程。随后删除缓存目录并重启」,如果切分点正好落在两句之间,那么「删除缓存目录」所在的 chunk 就丢掉了「先结束残留进程」这个前提,检索到它的用户会得到一个不完整、而且实际会复发的答案。80 字符的重叠让相邻 chunk 各自都保有对方的尾巴。

def chunk(blocks, size=500, overlap=80):
    chunks, current = [], ''
    for block in blocks:
        text = f"{' > '.join(block.heading_path)}\n{block.text}"
        if current and len(current) + len(text) > size:
            chunks.append(current)
            current = current[-overlap:]     # 尾部带进下一个 chunk
        current += text + '\n'
    if current.strip():
        chunks.append(current)
    return chunks

这个实现有意保持粗糙,有两个地方值得注意。第一,标题路径被拼进了每个 chunk 的正文,所以它会占用 size 预算——层级很深的文档,实际正文可能只剩三百多字符。第二,current[-overlap:] 是按字符切的,会从半句话中间开始,这在中文里可以接受,因为中文没有词边界问题。

粒度怎么选不该拍脑袋。300 / 500 / 800 三档的实际差异是第 11 篇 OFAT 实验的第一组对比,这里只给默认值。

入库时每个 chunk 要记全五个字段:

{'document_id': ..., 'chunk_index': 12, 'heading_path': ['第四章 故障排查', '4.3 启动异常'],
 'page': 47, 'source_offset': 8213}

chunk_index 用于把相邻 chunk 拼回上下文,pagesource_offset 是第 12 篇引用定位的锚点。这些字段现在不写,到第 12 篇就得重跑一遍全量索引。

异步任务:Redis 和 arq 到这里才有必要

前七篇没有引入队列,因为没有需要排队的东西。索引任务是第一个真实需求:耗时几十秒、可以失败重试、不需要用户等着。

async def enqueue_document(document_id: UUID, file_key: str):
    await documents.mark(document_id, 'queued')
    await redis.enqueue_job('index_document', str(document_id), file_key,
                            _job_id=f'index:{document_id}')

async def index_document(ctx, document_id: str, file_key: str):
    await documents.mark(document_id, 'parsing', progress=0.1)
    blocks = await parse(file_key)
    chunks = chunk(blocks)
    await documents.mark(document_id, 'embedding', progress=0.4,
                         total_chunks=len(chunks))
    vectors = await embed_in_batches([c.text for c in chunks], batch_size=64)
    await chunks_repo.replace(document_id, chunks, vectors)
    await documents.mark(document_id, 'ready', progress=1.0)

_job_id 是去重键。用户连点两次上传按钮,或者前端重试了一次请求,arq 会因为 job id 重复而丢弃后来的那次入队,不会跑出两个并发的索引任务。

chunks_repo.replace 而不是 insert:在一个事务里先删掉该文档已有的 chunk,再插入新的。这让重建索引变成幂等操作——第 7 篇的结论在这里复用。

状态机比一般的任务队列多两个字段:

alter table documents
  add column status text not null default 'queued'
    check (status in ('queued','parsing','embedding','ready','failed','deleting','deleted')),
  add column progress real not null default 0,
  add column error_code text,
  add column retry_count int not null default 0;

error_code 存结构化错误码,不存堆栈——理由和第 7 篇的审计表一样。retry_count 超过三次就停在 failed,不再自动重试。这时需要人看一眼:parse_failed 可能是加密 PDF,embedding_rate_limited 是配额问题,两者的处理方式完全不同。管理员在界面上点「重建索引」重新入队。

自动重试也要分类:embedding_rate_limitedupstream_timeout 可以退避重试;parse_failedunsupported_format 重试一百次也是同样的结果,直接进 failed

状态接口:前端进度条的数据源

@router.get('/documents/{document_id}/status')
async def status(document_id: UUID, user: User = Depends(current_user)):
    doc = await documents.visible_to(document_id, user)
    if not doc:
        raise HTTPException(404, 'document_not_found')
    return {
        'status': doc.status, 'progress': doc.progress,
        'processed_chunks': doc.processed_chunks,
        'total_chunks': doc.total_chunks, 'error': doc.error_code,
    }

documents.visible_to 而不是 documents.get:文档 ID 是 UUID,但不能因为它难猜就跳过权限校验。这里返回 404 而不是 403,是为了不泄漏「这个 ID 存在」这个信息。第 6 篇的可见性规则在这里第一次作用到文档上。

Preview
202 与状态轮询的完整时序

图里浏览器和 worker 之间没有任何直接连线。它们唯一的通信媒介是数据库里那行状态。这带来的好处是:用户关掉页面、刷新、甚至换一台设备重新登录,索引都照样跑完,重新打开时进度条能接上。

轮询间隔取 1 秒。索引是分钟级任务,200ms 的轮询只是在给自己的数据库做压测。前端拿到 ready 停止,拿到 failed 展示重试按钮——这个接口的完整消费方在第 14 篇。

删除:三个状态而不是一个 DELETE

删除文档要动两个存储:数据库里的 chunk 和对象存储里的原文件。它们无法在同一个事务里提交,所以用状态机把中间态显式表达出来:

async def delete_document(document_id: UUID, user: User):
    await documents.mark(document_id, 'deleting')       # 1. 先标记,检索立刻看不到
    await chunks_repo.delete_by_document(document_id)   # 2. 事务内删 chunk
    await object_store.delete(doc.file_key)             # 3. 删文件,失败可重试
    await documents.mark(document_id, 'deleted')        # 4. 终态

先标记 deleting 的意义在于:第 2 步和第 3 步之间如果进程崩了,文档已经从检索结果里消失了,不会有用户看到指向已删文件的引用。剩下的垃圾由清理任务扫 deleting 状态收尾。

这就是为什么检索查询必须带 documents.status = 'ready'。少这一个条件,正在索引的半成品 chunk 和正在删除的残留 chunk 都会被召回。这个条件在第 11 篇的两路 SQL 里都会出现。

测试与验收

docker compose up -d postgres redis
uv run pytest tests/documents/test_parsers.py tests/documents/test_chunking.py -q
uv run pytest tests/documents/test_index_job.py -q

四种格式各准备一个 fixture,PDF 那份要故意包含一个跨页的段落——它能测出解析器有没有把页脚当正文。

验收标准:上传接口在 2 秒内返回 202;重复上传因去重键只跑一个任务;索引进度可查询且单调递增;任务中途失败后重试不留下重复 chunk;不可重试的错误直接进 failed 并带 error_codepageheading_path 落库完整;检索查询过滤掉非 ready 文档。

现在库里有可检索的文档向量了。下一篇把第 8 篇导出的群聊接上来——你会看到这一篇的固定窗口切分在它身上完全失效,以及两个来源并进同一张表之后,说法不一致时该怎么办。