尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

AI Agent文档层设计:DocuQueue如何构建可检索的RAG异步管道

AI Agent文档层设计:DocuQueue如何构建可检索的RAG异步管道 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 有多个 SheetHTML 里混杂脚本和样式。如果让 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, encodingutf-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(funsupported 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_size800上下文更完整但检索粒度变粗语义更聚焦但信息容易碎片化chunk_overlap100减少上下文割裂但增加向量存储节省存储但可能切断语义中文场景下字符数并不等于 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( fno handler for {task.get(task_type)} ) await handler(task) finally: self._queue.task_done()worker 的职责是从队列取任务、按类型分发、调用处理函数。生产环境需要把MemoryTaskQueue替换成 Redis Streams 或 RabbitMQ并且任务执行成功后要主动确认执行失败时决定重试还是进入死信。这里有一个容易被忽略的点重试不应该简单地把整个任务重新入队。更合理的做法是记录任务已经推进到哪个阶段比如“解析完成但切块失败”重试时直接从切块阶段继续避免重复解析大文件。3.5 向量化与向量入库文本切块之后需要经过嵌入模型变成向量。示例中使用 SentenceTransformerfrom 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(urlhttp://localhost:6333) COLLECTION_NAME documents def create_collection(vector_size: int) - None: client.recreate_collection( collection_nameCOLLECTION_NAME, vectors_configVectorParams( sizevector_size, distanceDistance.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( idpoint_id, vectorvector, 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_nameCOLLECTION_NAME, pointspoints, )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( keyteam, matchmodels.MatchValue(valueteam), ) for team in team_filter ] ) return client.search( collection_nameCOLLECTION_NAME, query_vectorquery_vector, limittop_k, query_filterquery_filter, ) app FastAPI() app.post(/query, response_modelQueryResponse) def query_documents(req: QueryRequest): vector embed(req.query) hits search_chunks( query_vectorvector, team_filterreq.user_teams or None, top_kreq.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 [], }, timeout10, ) 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任务处理查询任务状态状态变为 succeededchunk 入库在向量库中查询 doc_id存在多条 chunk 记录检索返回调用 /query返回相关文本score 符合预期权限过滤用非 infra team 查询不返回 infra 文档失败重试临时停掉向量库再提交文档任务失败后按策略重试不重复入库队列监控查看队列深度指标没有无限积压资源占用连续提交多个大文件内存和 CPU 处于可控范围
返回列表