AI Agent 能调模型、能写代码、能操作工具,但一遇到企业内部 PDF、Word、Excel 组成的长文档,问题就出来了:上下文窗口放不下,文件格式又杂,文档还在不断更新。这个问题不能靠 prompt 修补,需要一个专门的文档层(Document Layer)来承接文档的清洗、切分、向量化和检索。DocuQueue 的名字点明了它的设计方向:Document 加 Queue,让 Agent 不再和原始文件逻辑纠缠,而是从文档层拿到结构化、可检索、有权限边界的结果。下面把 DocuQueue 当作一类文档层组件的统称来讨论,先梳理它解决什么问题,再落到最小实现、Agent 接入、运行验证、常见坑和生产落地建议。
1. 为什么 AI Agent 需要独立的文档层
1.1 Agent 直接处理文档的痛点
Agent 在处理少量文本时,直接拼装进 prompt 问题不大。但真实业务场景通常是几百份合同、产品手册、内部 Wiki、客服记录,还有不断新增的 Markdown 和网页快照。把这些原始文件直接交给模型,上下文窗口首先就会超限,紧随而来的是 token 费用、响应延迟和输出不稳定。更麻烦的是,每次会话都要重复处理同一份文档,缺少一次构建、多次复用的机制。
文件格式的多样性比想象中更复杂。PDF 可能是文字版,也可能是扫描版;Word 文档里有页眉页脚、表格、批注;Excel 有多个 Sheet;HTML 里混杂脚本和样式。如果让 Agent 自己写解析逻辑,同一份文档在不同任务里可能解析结果不一致,甚至有的模型工具只能看到文本,看不到表格结构。文档层要做的第一件事,就是把“文件读取”和“内容理解”解耦。
文档不是静态的。合同会更新,知识库会删除旧版本,新资料发布后需要在几分钟内生效。如果每次都由 Agent 触发解析,既会产生延迟,也会造成大量重复计算。权限问题同样无法回避:不同用户、不同团队应该只能检索到自己有权限访问的文档。这些规则如果散落在 Agent 调用链里,很难统一治理。
所以更合理的分工是:Agent 只负责基于检索结果进行推理,不负责把文件变成结果。文档的接入、解析、切块、向量化、索引、更新、权限过滤和检索,都应该由一个独立组件承担。这个组件就是文档层。
1.2 文档层的定位与边界
文档层可以类比成业务应用和数据库之间的 ORM:上层不需要知道 SQL 怎么执行,只需要拿到数据。文档层面对 Agent 暴露的是“检索能力”,而不是“文件系统”。
| 文档层负责 | 文档层不负责 |
|---|---|
| 文档上传、格式识别、内容清洗 | Agent 对话策略和 prompt 拼装 |
| 文本切分、向量化、索引构建 | LLM 调用和模型微调 |
| 文档元数据管理、版本更新 | Agent 工具调度和任务编排 |
| 权限过滤、检索结果返回 | 最终答案生成 |
| 任务状态、重试、监控 | 前端交互和业务逻辑 |
边界清楚之后,文档层可以独立升级解析能力,Agent 侧不需要改动。反过来,Agent 想切换模型或调整 prompt,也不会影响已经建立的文档索引。
1.3 DocuQueue 要解决的核心问题
DocuQueue 这类组件的核心不只是“封装文件解析”,而是把文档处理建模成异步任务队列。文档处理是一条有状态、可重试、需要监控的管道,不是一次同步函数调用。
一次完整的文档处理可以简化为:
提交文档 -> 任务入队 -> 工作节点拉取 -> 解析清洗 -> 文本切块 -> 向量化 -> 写入向量库 -> 更新任务状态用队列来组织这条链路,至少有四个好处:
- 大文件解析和向量化耗时较长,不能让用户请求一直阻塞等待;
- 每一步都可能失败,任务需要有 pending、running、succeeded、failed 等明确状态;
- 失败后要支持有限次重试,重试只处理失败的阶段,不能每次从头开始;
- 不同来源的文档可以共享同一个处理管道,便于控制并发和观察积压。
这就是 DocuQueue 名字的由来:文档处理不是一次查询,而是一条可以被排队、调度和追踪的流水线。
2. 文档层的核心组成与数据模型
2.1 文档从进入到可检索经历哪些阶段
要把文档层设计清楚,先要把处理阶段拆开。每个阶段都有输入、输出和最容易失败的位置。
| 阶段 | 输入 | 输出 | 常见的失败点 |
|---|---|---|---|
| 接入 | 文件流、URL、原始文本 | 文档记录 | 格式不支持、上传中断 |
| 解析 | 文件路径 | 纯文本 | PDF 加密、扫描件无文字层 |
| 清洗 | 纯文本 | 规范文本 | 乱码、表格错乱、页眉页脚混入 |
| 切分 | 规范文本 | chunk 列表 | 切分粒度不合理,切断语义 |
| 嵌入 | chunk 文本 | 向量 | 模型加载失败、显存不足 |
| 索引 | 向量加 metadata | 向量索引 | 向量库连接失败、字段类型不匹配 |
| 就绪 | 索引 | 检索结果 | 权限过滤条件写错,返回无权访问内容 |
实际项目不一定每个阶段都独立成服务,但至少要在日志中保留阶段标识。否则文档检索结果不对时,很难判断是解析丢了内容,还是切分切坏了上下文,还是向量库没有更新。
2.2 文档项、文档块和任务对象
文档层至少维护三类核心对象:文档、文档块、处理任务。
文档对象记录“源文件是什么、属于谁、当前状态如何”:
{ "doc_id": "doc-001", "file_name": "server-setup.pdf", "content_type": "application/pdf", "owner": "user-1001", "team": "infra", "status": "succeeded", "created_at": "2025-01-01T10:00:00Z" }文档块对象用于检索,它才是真正会被 Agent 拿去做上下文的内容:
{ "chunk_id": "doc-001:00007", "doc_id": "doc-001", "chunk_index": 7, "text": "部署前需要先确认 Docker 版本不低于 24.0...", "token_count": 260, "metadata": { "team": "infra", "page": 3, "source_url": "https://wiki.internal/server-setup" } }处理任务对象描述管道执行到哪里了:
{ "task_id": "task-abc123", "doc_id": "doc-001", "task_type": "ingest", "status": "failed", "attempts": 2, "last_error": "PDF 文件已加密,无法提取文本", "updated_at": "2025-01-01T10:01:23Z" }chunk 的 metadata 非常关键。权限过滤、来源追踪、引用原文、按页面筛选,都依赖 metadata。不要等检索阶段再补字段,因为那时原始上下文可能已经丢失。
2.3 队列与状态机
文档任务的状态机可以做成这样:
pending -> running -> succeeded | v failed | v retry (重新进入 pending 或 running) | v dead letterfailed不一定是终点。允许有限次重试有助于临时故障恢复,比如向量库短时间不可用。但如果超过重试次数仍然失败,任务应该进入死信队列,由人工检查并决定是删除还是修复后重新入队。
队列选型可以按团队规模来:
| 方案 | 优点 | 适用的场景 |
|---|---|---|
| asyncio.Queue | 零依赖,代码简单 | 本地开发、单机演示 |
| Redis Streams | 持久化、可靠,消费组机制成熟 | 中小团队,文档量中等 |
| RabbitMQ | 路由灵活,ack 机制完善 | 已有消息中间件,需要复杂的路由策略 |
| Kafka | 高吞吐、可回放 | 大规模事件流,文档更新频繁 |
不建议一上来就上 Kafka。如果日均处理文档只有几百份,用 Redis 队列加一套可靠的重试机制已经足够,架构复杂度越低越好维护。
3. 用 Python 实现一个 DocuQueue 最小版本
3.1 环境准备与项目结构
下面的示例用于表达文档层的设计思路,不绑定某个具体开源仓库。落地到自己项目时,需要把包名、存储实现、模型路径替换成实际环境。
建议使用 Python 3.10 及以上版本。
mkdir docuqueue-demo cd docuqueue-demo python -m venv .venv source .venv/bin/activate pip install fastapi uvicorn pydantic根据实际需要,可能还会用到:
pip install pypdf python-docx sentence-transformers qdrant-client最小项目结构可以这样组织:
docuqueue-demo/ ├── app/ │ ├── main.py │ ├── models.py │ ├── parser.py │ ├── chunker.py │ ├── queue.py │ ├── embedding.py │ ├── vector_store.py │ └── retrieval.py └── tests/模块拆分原则是:一个模块只负责一个阶段。解析、切分、嵌入、索引分别独立,方便替换实现。比如把本地文件解析换成对象存储流式读取时,只需要改 parser 模块。
3.2 定义核心数据对象
先定义文档、任务、状态枚举:
from enum import Enum from typing import Optional from datetime import datetime from pydantic import BaseModel class TaskStatus(str, Enum): PENDING = "pending" RUNNING = "running" SUCCEEDED = "succeeded" FAILED = "failed" CANCELED = "canceled" class Document(BaseModel): doc_id: str file_name: str content_type: str = "text/plain" owner: str = "default" team: str = "default" status: TaskStatus = TaskStatus.PENDING created_at: datetime = datetime.utcnow() class DocumentChunk(BaseModel): chunk_id: str doc_id: str chunk_index: int text: str token_count: int = 0 metadata: dict = {} class ProcessingTask(BaseModel): task_id: str doc_id: str task_type: str = "ingest" status: TaskStatus = TaskStatus.PENDING attempts: int = 0 last_error: Optional[str] = None用 Pydantic 或类似库定义数据对象,能减少字段拼写错误。status使用枚举而不是裸字符串,可以避免状态值写错后无法被程序及时发现。
3.3 文档解析与切块
解析模块负责把文件变成纯文本。不同格式的分支可以后续扩展:
def extract_text(file_path: str, content_type: str) -> str: if content_type == "text/plain": with open(file_path, "r", encoding="utf-8") as f: return f.read() if content_type == "application/pdf": from pypdf import PdfReader reader = PdfReader(file_path) return "\n".join( page.extract_text() or "" for page in reader.pages ) raise ValueError(f"unsupported content type: {content_type}")切块模块使用固定大小加重叠的方式,这是最常见的起步方案:
def split_text( text: str, chunk_size: int = 800, chunk_overlap: int = 100, ) -> list[str]: if chunk_size <= chunk_overlap: raise ValueError("chunk_size 必须大于 chunk_overlap") chunks = [] start = 0 text_len = len(text) while start < text_len: end = min(start + chunk_size, text_len) chunks.append(text[start:end]) if end >= text_len: break start = end - chunk_overlap return chunks切块参数是文档层里最影响检索效果的部分之一。
| 参数 | 默认值 | 调大 | 调小 |
|---|---|---|---|
| chunk_size | 800 | 上下文更完整,但检索粒度变粗 | 语义更聚焦,但信息容易碎片化 |
| chunk_overlap | 100 | 减少上下文割裂,但增加向量存储 | 节省存储,但可能切断语义 |
中文场景下字符数并不等于 token 数。切块后最好通过分词器或模型 tokenizer 计算token_count,它是后期排查上下文超限和费用分析的重要指标。
3.4 队列与任务调度
最小演示可以用asyncio.Queue实现一个内存队列:
import asyncio from typing import Callable, Awaitable class MemoryTaskQueue: def __init__(self): self._queue = asyncio.Queue() self._handlers = {} def register( self, task_type: str, handler: Callable[[dict], Awaitable[None]], ) -> None: self._handlers[task_type] = handler async def submit(self, task: dict) -> None: await self._queue.put(task) async def run_worker(self) -> None: while True: task = await self._queue.get() try: handler = self._handlers.get(task.get("task_type")) if handler is None: raise ValueError( f"no handler for {task.get('task_type')}" ) await handler(task) finally: self._queue.task_done()worker 的职责是从队列取任务、按类型分发、调用处理函数。生产环境需要把MemoryTaskQueue替换成 Redis Streams 或 RabbitMQ,并且任务执行成功后要主动确认,执行失败时决定重试还是进入死信。
这里有一个容易被忽略的点:重试不应该简单地把整个任务重新入队。更合理的做法是记录任务已经推进到哪个阶段,比如“解析完成,但切块失败”,重试时直接从切块阶段继续,避免重复解析大文件。
3.5 向量化与向量入库
文本切块之后,需要经过嵌入模型变成向量。示例中使用 SentenceTransformer:
from sentence_transformers import SentenceTransformer _model = SentenceTransformer( "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2" ) def embed(text: str) -> list[float]: return _model.encode(text).tolist()向量入库时,要同时写入 metadata,否则后续无法做权限过滤:
from qdrant_client import QdrantClient from qdrant_client.models import Distance, VectorParams, PointStruct client = QdrantClient(url="http://localhost:6333") COLLECTION_NAME = "documents" def create_collection(vector_size: int) -> None: client.recreate_collection( collection_name=COLLECTION_NAME, vectors_config=VectorParams( size=vector_size, distance=Distance.COSINE, ), ) def upsert_chunks( chunks: list[dict], vectors: list[list[float]], ) -> None: points = [] for chunk, vector in zip(chunks, vectors): point_id = f"{chunk['doc_id']}:{chunk['chunk_index']}" # 注意:point_id 类型以所用向量库实际支持为准 points.append( PointStruct( id=point_id, vector=vector, payload={ "doc_id": chunk["doc_id"], "chunk_index": chunk["chunk_index"], "team": chunk["metadata"].get("team", ""), "owner": chunk["metadata"].get("owner", ""), "text": chunk["text"], }, ) ) client.upsert( collection_name=COLLECTION_NAME, points=points, )metadata 里至少要包含doc_id、chunk_index、team、owner和原文。doc_id + chunk_index应该作为幂等键,避免重试时重复写入同一块内容。
3.6 提供检索接口
文档层对外可以暴露一个/query接口:
from fastapi import FastAPI from pydantic import BaseModel from qdrant_client import models class QueryRequest(BaseModel): query: str top_k: int = 5 user_teams: list[str] = [] class QueryResponse(BaseModel): chunks: list[dict] def search_chunks( query_vector: list[float], team_filter: list[str] | None = None, top_k: int = 5, ): query_filter = None if team_filter: query_filter = models.Filter( should=[ models.FieldCondition( key="team", match=models.MatchValue(value=team), ) for team in team_filter ] ) return client.search( collection_name=COLLECTION_NAME, query_vector=query_vector, limit=top_k, query_filter=query_filter, ) app = FastAPI() @app.post("/query", response_model=QueryResponse) def query_documents(req: QueryRequest): vector = embed(req.query) hits = search_chunks( query_vector=vector, team_filter=req.user_teams or None, top_k=req.top_k, ) return QueryResponse( chunks=[ { "doc_id": hit.payload["doc_id"], "chunk_index": hit.payload["chunk_index"], "text": hit.payload["text"], "score": hit.score, "team": hit.payload["team"], } for hit in hits ] )这个接口的关键点在于:权限过滤在向量检索阶段完成,而不是等结果返回后由应用层丢弃。因为向量检索返回时已经把内容带到内存里,后置过滤无法真正防止越权请求。
4. 接入 AI Agent:从文档层到 RAG 检索
4.1 Agent 调用文档层的两种模式
Agent 接入文档层,最常用的是同步查询模式:Agent 收到问题后,先调用文档层检索接口获取相关上下文,再把上下文和用户问题一起交给 LLM 生成回答。
另一种是事件订阅模式。文档层在文档完成索引、更新或删除时发送事件,Agent 或周边系统根据事件更新自己的缓存、通知用户或触发后续工作流。
| 模式 | 调用方向 | 优点 | 适用场景 |
|---|---|---|---|
| 同步查询 | Agent 调用 /query | 实时性好、实现简单 | 对话 RAG、问答助手 |
| 事件订阅 | 文档层推送事件 | 异步、减少轮询 | 知识库监控、文档更新提醒 |
两种模式可以组合。Agent 回答问题时同步检索,同时后台订阅文档更新事件,用于刷新热门文档缓存。
4.2 检索接口设计
检索接口的请求和返回最好保持通用,不要在接口里绑定某个具体 Agent 框架。
{ "query": "如何配置日志采集", "top_k": 5, "user_teams": ["infra"], "score_threshold": 0.35 }返回结构:
{ "chunks": [ { "doc_id": "doc-023", "chunk_index": 3, "text": "日志采集需要先安装 agent,然后配置采集路径...", "score": 0.78, "metadata": { "team": "infra", "source_url": "https://wiki.internal/log-agent" } } ] }top_k决定最多返回多少个块。score_threshold用于过滤相似度过低的噪声结果。实际项目中,阈值需要根据测试集调整,不要直接使用一个拍脑袋的数字。
4.3 权限过滤与元数据控制
权限过滤是文档层最容易被轻视的部分。写入时没有把team、owner放进 metadata,检索时就无法过滤。即使后来补上,历史文档也要重跑索引。
一个简化的权限规则可以是:
- 文档层保存文档的可见团队列表;
- 每个 chunk 的 metadata 继承文档的团队信息;
- Agent 发起检索时携带当前用户所属团队;
- 向量检索强制添加 team 过滤条件;
- 对返回结果再校验一次文档可访问范围。
不要把权限判断全部交给 Agent 代码。Agent 代码可能被替换,也可能在处理多轮对话时遗漏过滤条件。文档层作为数据出口,必须自己守住边界。
4.4 与 Agent 框架集成示例
文档层接口是 HTTP 接口,所以无论 LangChain、LlamaIndex 还是自研 Agent,都可以封装成一个工具或检索器。
下面是一个通用的 Python 检索器封装:
import requests class DocumentLayerRetriever: def __init__(self, endpoint: str, top_k: int = 5): self.endpoint = endpoint self.top_k = top_k def get_context( self, query: str, user_teams: list[str] | None = None, ) -> list[str]: response = requests.post( f"{self.endpoint}/query", json={ "query": query, "top_k": self.top_k, "user_teams": user_teams or [], }, timeout=10, ) response.raise_for_status() return [ { "text": item["text"], "score": item["score"], "source": item["doc_id"], } for item in response.json()["chunks"] ]接入 Agent 时,只需把get_context的返回值拼进系统提示词或作为工具参数。要注意给 HTTP 请求设置超时并处理超时分支,不能因为文档层临时不可用而让整个 Agent 挂起。
5. 运行验证与可观测性
5.1 验证文档入库链路
假设 FastAPI 服务运行在 8000 端口,可以按下面顺序验证。
启动服务:
uvicorn app.main:app --host 0.0.0.0 --port 8000提交一份测试文档:
curl -X POST http://localhost:8000/documents \ -H "Content-Type: application/json" \ -d '{ "doc_id": "doc-001", "file_name": "guide.txt", "content_type": "text/plain", "owner": "user-1001", "team": "infra" }'查询任务状态:
curl http://localhost:8000/tasks/task_doc-001调用检索接口:
curl -X POST http://localhost:8000/query \ -H "Content-Type: application/json" \ -d '{ "query": "如何配置环境", "top_k": 3, "user_teams": ["infra"] }'预期结果:
- 提交文档后返回任务 ID,任务状态为
pending; - 等待异步处理完成后,任务状态变为
succeeded; - 向量库中能看到该文档对应 chunk 记录;
/query返回的文本与测试文档内容相关,且带score和 metadata。
如果任务长期停留在pending,说明 worker 没有消费队列,要检查 worker 是否启动、队列连接是否正常。
5.2 可观测性:日志、指标和追踪字段
文档层的排错不能只靠“有没有报错”。要把每一次任务处理过程记录下来,方便回溯。
| 类型 | 关键字段 | 用途 |
|---|---|---|
| 日志 | task_id、doc_id、阶段、status、duration_ms、error | 定位单个任务失败原因 |
| 指标 | queue_depth、processed_total、chunk_total、embed_latency | 判断系统容量和瓶颈 |
| 追踪 | trace_id、span_id、阶段名 | 串联从提交文档到检索的完整链路 |
日志建议使用结构化 JSON 格式。例如:
{ "time": "2025-01-01T10:01:23Z", "level": "error", "logger": "docuqueue.parser", "task_id": "task-abc123", "doc_id": "doc-001", "stage": "parse", "status": "failed", "duration_ms": 320, "error": "PDF 文件已加密" }有了结构化日志,后续查问题时可以直接在日志平台按task_id定位,而不是翻文件搜关键字。
5.3 发布前检查清单
这里给出一份可复制的检查清单:
| 检查项 | 检查方式 | 预期结果 |
|---|---|---|
| 文档上传 | curl 提交测试文档 | 返回 task_id,状态 pending |
| 任务处理 | 查询任务状态 | 状态变为 succeeded |
| chunk 入库 | 在向量库中查询 doc_id | 存在多条 chunk 记录 |
| 检索返回 | 调用 /query | 返回相关文本,score 符合预期 |
| 权限过滤 | 用非 infra team 查询 | 不返回 infra 文档 |
| 失败重试 | 临时停掉向量库再提交文档 | 任务失败后按策略重试,不重复入库 |
| 队列监控 | 查看队列深度指标 | 没有无限积压 |
| 资源占用 | 连续提交多个大文件 | 内存和 CPU 处于可控范围 |