AI Agent 的文档处理能力,正在从“能读文件”升级为“会管理知识”。但很多团队搭建 Agent 应用时,都会卡在同一个环节:文档接进来之后,怎么存、怎么切、怎么检索,才能让模型真正用得上?DocuQueue 给出了一种思路——在 Agent 与文档之间增加一个独立的 Document Layer(文档层),把文档从“静态资源”变成“可调度的数据流”。
这篇文章会从文档层的概念讲起,拆解 DocuQueue 的架构设计和核心能力,并提供一个可运行的轻量级 Python 实现。你可以把它当成一个简化版参考,结合自己的业务场景去改造和扩展。
1. 为什么 AI Agent 需要一个文档层
1.1 直接读文件为什么不够用
早期的 Agent 处理文档,最常见的方式是:把 PDF、Word、Markdown 一次性读进来,拼接成一段很长的文本,塞进 Prompt。这种方式在演示环境里效果尚可,但进入真实业务场景后问题很明显。
Token 限制是第一个瓶颈。一个几十页的 PDF,转成文本后可能超过几万字,直接全部塞给模型,很快就把上下文窗口占满。更麻烦的是,模型面对大量无关内容时,注意力会被稀释,回答质量大幅下降。文档更新频率是第二个问题。如果 Agent 每次读取的都是最新的完整文档,体积会越来越大,成本也会不断上涨。第三个问题是格式差异。PDF 需要解析,Word 需要转换,HTML 需要清洗,每种格式都有不同的处理逻辑,如果这些逻辑全部耦合在 Agent 主流程里,代码会迅速变得混乱。
所以,我们需要在 Agent 和原始文档之间加一层“中间层”,专门负责处理文档的接入、转换、存储、检索和更新。这层基础设施就是 Document Layer。
1.2 DocuQueue 解决什么问题
DocuQueue 本质上是一个针对 AI Agent 场景设计的文档队列系统。它借鉴了消息队列的思想,把文档当作消息来处理。每一份文档进入系统后,会经历“入队 → 解析 → 切分 → 增强 → 索引 → 可检索”的完整生命周期。Agent 不再直接接触原始文件,而是通过 Document Layer 提供的统一接口获取所需内容。
这个设计带来的好处非常直接:
- 文档处理与 Agent 执行解耦,Agent 不用关心文件从哪来、是什么格式。
- 文档更新可以异步完成,新版本入库后,Agent 自动检索到最新内容。
- 检索接口统一,无论底层是向量库、Elasticsearch 还是纯内存索引,对 Agent 来说没有差异。
1.3 Document Layer 在 Agent 架构中的位置
在一个典型的 RAG(检索增强生成)架构中,Document Layer 位于“数据源”和“Agent/模型”之间。你可以把它理解为三明治的中间层。
整个数据流向可以概括为:
- 原始文档从各种渠道进入系统(上传、爬取、消息队列、定时任务拉取)。
- Document Layer 对文档进行解析和标准化,转成统一的 Document 对象。
- 文档对象经过切分、清洗、元数据标记后被写入存储引擎。
- Agent 提出请求时,Document Layer 执行检索,返回最相关的文档片段。
- Agent 将检索结果与用户问题组合,交给大模型生成最终回答。
没有这个中间层,Agent 的代码里会堆满文件解析、格式判断、Token 截断之类的脏活累活。有了它,Agent 只需要面向接口编程。
2. DocuQueue 核心设计思路
2.1 文档生命周期:从入队到可检索
DocuQueue 把文档的生命周期划分为五个阶段,每个阶段都可以独立扩展。理解这个生命周期,是理解整个框架的关键。
第一阶段是接入。不同来源的文档被统一封装成 job 进入队列,这个阶段不关心文档内容,只负责收拢数据源。第二阶段是解析。系统根据文件类型选择合适的解析器,把 PDF、Word、Markdown 转成纯文本。第三阶段是切分。长文本按照一定的策略被切成多个 chunk,每个 chunk 是后续检索的最小单位。第四阶段是增强。系统为每个 chunk 补充元数据,包括文档编号、章节标题、关键词、时间戳等。第五阶段是索引与存储。chunk 被写入检索系统,Agent 可以通过查询接口获取。
整个生命周期中的每一步都支持异步执行。文档进入系统后,Agent 不需要阻塞等待处理完成,而是可以通过状态查询接口获取处理进度。
2.2 队列模型:为什么用队列而不是直连
在 DocuQueue 的设计中,队列是核心抽象,它的存在有三个层面的原因。
第一个原因是削峰填谷。当大量文档同时涌入时,如果 Agent 同步等待全部处理完成,响应时间会非常不稳定。通过队列缓冲,系统可以按固定速率处理文档,避免资源被打满。
第二个原因是故障隔离。队列使得文档处理可以重试。如果某份文档解析失败,系统会保存任务状态,并在一段时间后重新尝试,而不是直接崩溃或返回错误给用户。
第三个原因是处理顺序控制。不同来源的文档可能有不同的优先级。比如用户手动上传的文档需要立即处理,而定时爬取的网页可以放在低优先级队列中慢慢消费。这种能力在同步直连模型里很难实现。
2.3 可插拔设计:存储与检索的抽象
DocuQueue 在设计上允许替换底层存储和检索组件。在开发环境中,你可以使用内存存储;数据量小的时候,可以用 SQLite 或 JSON 文件;数据量增长后,切换到向量数据库或 Elasticsearch。
这种可插拔设计带来的工程价值在于:团队可以先快速验证产品逻辑,再根据实际数据规模和查询需求调整基础设施。Demo 阶段和上线阶段用同一套上层 API,迁移成本是可控的。
3. 环境准备与项目结构
3.1 技术选型说明
下面我们将用 Python 从零实现一个简化版 Document Layer。选择 Python 的原因是它在 AI Agent 生态中支持最完善,且示例代码便于阅读和修改。
建议环境版本:
- Python 3.10+
- 依赖库:pydantic(数据模型校验)
- 可选依赖:openai(Agent 接入)、redis(队列存储)、numpy(简单的向量计算)
本文的示例重点是架构思路和核心流程,不会依赖任何重量级框架。你需要根据项目实际情况调整版本,重点是理解设计模式。
3.2 项目目录结构
首先建立一个清晰的项目结构。清晰的目录是良好架构的第一步。
docuqueue/ ├── __init__.py ├── models.py # 数据模型定义 ├── queue.py # 文档队列 ├── chunker.py # 文档切分器 ├── parser.py # 文档解析器 ├── storage.py # 存储与检索接口 ├── docuqueue.py # 核心门面类 └── examples/ └── demo_agent.py # Agent 接入示例目录划分的职责非常明确:
- models.py 定义统一的数据结构
- parser.py 负责把原始文件转为文本
- chunker.py 负责把长文本切成片段
- storage.py 负责 chunk 的存储和检索
- queue.py 管理文档处理任务的有序执行
- docuqueue.py 对外提供面向 Agent 的 API
4. 从零实现一个轻量 Document Layer
下面我们实现一个简化版 DocuQueue。代码聚焦核心流程,方便你在自己的机器上运行和调试。
4.1 定义核心数据模型
数据模型是整个系统的基础。我们定义三个关键类:Document 代表一份文档,Chunk 代表切分后的文本片段,DocStatus 描述文档的处理状态。
# 文件路径:docuqueue/models.py from __future__ import annotations import time import uuid from enum import Enum from typing import Any, Optional from pydantic import BaseModel, Field class DocStatus(str, Enum): PENDING = "pending" PARSING = "parsing" CHUNKING = "chunking" INDEXED = "indexed" FAILED = "failed" class Document(BaseModel): """统一文档模型,所有来源的文档都会被转换为该结构。""" doc_id: str = Field(default_factory=lambda: uuid.uuid4().hex) source_name: str = Field(..., description="文档来源名称,例如文件名或URL") content: str = Field(default="", description="解析后的纯文本内容") metadata: dict[str, Any] = Field(default_factory=dict, description="文档元数据") status: DocStatus = DocStatus.PENDING created_at: float = Field(default_factory=time.time) updated_at: float = Field(default_factory=time.time) def update_status(self, status: DocStatus) -> None: self.status = status self.updated_at = time.time() class Chunk(BaseModel): """检索单元,是文档切分后的最小片段。""" chunk_id: str = Field(default_factory=lambda: uuid.uuid4().hex) doc_id: str = Field(...) index: int = Field(..., description="在文档中的序号,从0开始") content: str = Field(...) metadata: dict[str, Any] = Field(default_factory=dict) @property def position(self) -> str: """返回便于查看的位置描述""" return f"{self.doc_id}#{self.index}"代码里的关键设计点:
- 使用 pydantic 的 BaseModel 做数据校验,确保字段类型正确。
- Document 和 Chunk 都有独立的 id,方便在队列、存储中追踪。
- metadata 字段是开放式的,允许不同来源的文档附带不同的附加属性。
4.2 实现文档解析器
解析器负责把各种格式的文件转为纯文本。真正的生产系统可能需要接入 PDF、Word、HTML 解析库。这里我们简化处理,同时演示“按扩展名分发解析器”的设计思路。
# 文件路径:docuqueue/parser.py from __future__ import annotations from pathlib import Path from typing import Callable from docuqueue.models import Document class DocumentParser: """文档解析器,根据文件扩展名选择对应解析逻辑。""" def __init__(self) -> None: # 注册解析器表:扩展名 -> 处理函数 self._parsers: dict[str, Callable[[Path], str]] = { ".txt": self._parse_txt, ".md": self._parse_markdown, ".json": self._parse_json, } def register_parser(self, extension: str, func: Callable[[Path], str]) -> None: """注册新的自定义解析器。""" self._parsers[extension] = func def parse(self, file_path: str | Path) -> Document: path = Path(file_path) if not path.exists(): raise FileNotFoundError(f"文件不存在: {path}") ext = path.suffix.lower() if ext not in self._parsers: # 不支持的文件类型按纯文本读取,避免直接丢弃 content = path.read_text(encoding="utf-8", errors="ignore") else: content = self._parsers[ext](path) doc = Document( source_name=path.name, content=content, metadata={ "file_path": str(path), "extension": ext, "file_size": path.stat().st_size, }, ) return doc def _parse_txt(self, path: Path) -> str: return path.read_text(encoding="utf-8", errors="ignore") def _parse_markdown(self, path: Path) -> str: # 真实项目中可以使用 markdown 库并保留标题结构 # 这里先按文本读取,读者可根据需要自行扩展 return path.read_text(encoding="utf-8", errors="ignore") def _parse_json(self, path: Path) -> str: # 把 JSON 转成紧凑的文本,方便后续切分与检索 import json with path.open("r", encoding="utf-8") as fp: data = json.load(fp) return json.dumps(data, ensure_ascii=False)这里体现了“可插拔”的思想。如果你需要支持 PDF,只需要安装一个 PDF 解析库,然后在初始化时调用 register_parser(".pdf", my_pdf_parser),无需修改其他代码。
4.3 实现文档切分器
切分是整个 Document Layer 中最核心的环节之一。切分策略直接影响检索质量。我们实现一个递归字符切分器,核心思路是:优先按段落切,段落太长时按句子切,句子太长时按固定长度切。
# 文件路径:docuqueue/chunker.py from __future__ import annotations from typing import Iterable from docuqueue.models import Chunk, Document class RecursiveCharacterChunker: """ 递归字符串切分器。 切分策略: 1. 优先使用段落分割符(\n\n) 2. 段落过长时使用句子分割符(。!?\n) 3. 仍然过长时按窗口大小截断 这种策略能最大限度保留语义完整性。 """ def __init__( self, chunk_size: int = 1000, chunk_overlap: int = 100, ) -> None: if chunk_size <= 0: raise ValueError("chunk_size 必须大于0") if chunk_overlap < 0 or chunk_overlap >= chunk_size: raise ValueError("chunk_overlap 必须在 [0, chunk_size) 范围内") self.chunk_size = chunk_size self.chunk_overlap = chunk_overlap def split_document(self, doc: Document) -> list[Chunk]: raw_text = doc.content.strip() if not raw_text: return [] sections = self._recursive_split(raw_text) chunks = [] for i, section in enumerate(sections): chunk = Chunk( doc_id=doc.doc_id, index=i, content=section, metadata={ "source_name": doc.source_name, "chunk_size": len(section), }, ) chunks.append(chunk) return chunks def _recursive_split(self, text: str) -> list[str]: # 如果文本本身很短,直接返回 if len(text) <= self.chunk_size: return [text] # 第一轮:按段落切 paragraphs = [p.strip() for p in text.split("\n\n") if p.strip()] if len(paragraphs) > 1: return self._merge_paragraphs(paragraphs) # 第二轮:按句子切 sentences = [s.strip() for s in text.split("。") if s.strip()] if len(sentences) > 1: return self._merge_sentences(sentences) # 第三轮:直接按窗口切 return [text[i : i + self.chunk_size] for i in range(0, len(text), self.chunk_size)] def _merge_paragraphs(self, paragraphs: list[str]) -> list[str]: result = [] buffer = "" for p in paragraphs: if len(buffer) + len(p) <= self.chunk_size: buffer = f"{buffer}\n\n{p}".strip() else: if buffer: result.append(buffer) buffer = p if buffer: result.append(buffer) return result def _merge_sentences(self, sentences: list[str]) -> list[str]: result = [] buffer = "" for s in sentences: candidate = f"{buffer}{s}" if buffer else s if len(candidate) <= self.chunk_size: buffer = candidate else: if buffer: result.append(buffer) buffer = s if buffer: result.append(buffer) return result说明几个细节:
- chunk_overlap 在简化实现中没有展示完整效果,但参数保留。真实实现中,相邻 chunk 需要保留重叠部分,以保证跨片段语义不丢失。
- 段落优先的切分策略适合文档结构明显(标题、章节)的场景。
- 句子优先的切分策略适合连续叙述型的文本。
4.4 实现队列与任务调度
队列负责管理文档处理的生命周期。这里我们不引入外部消息中间件,而是用 Python 内置的数据结构模拟队列语义,重点展示流程编排。
# 文件路径:docuqueue/queue.py from __future__ import annotations import time from collections import deque from typing import Callable from docuqueue.models import DocStatus, Document class DocumentQueue: """ 文档队列:管理文档处理任务的有序执行。 生产环境可以替换为 Redis Streams 或 RabbitMQ。 本示例专注演示队列语义。 """ def __init__(self) -> None: self._pending: deque[str] = deque() self._documents: dict[str, Document] = {} self._status_handlers: dict[DocStatus, Callable[[str], None]] = {} def enqueue(self, doc: Document) -> None: self._documents[doc.doc_id] = doc doc.update_status(DocStatus.PENDING) self._pending.append(doc.doc_id) print(f"[queue] 文档入队: {doc.source_name} ({doc.doc_id[:8]})") def register_handler(self, status: DocStatus, handler: Callable[[str], None]) -> None: """注册状态处理器,例如切分成功后的索引回调。""" self._status_handlers[status] = handler def process_next(self, chunker=None, storage=None) -> bool: """ 处理队列中的下一个文档。 如果队列为空返回 False。 """ if not self._pending: return False doc_id = self._pending.popleft() doc = self._documents[doc_id] try: # 状态流转:pending -> parsing -> chunking -> indexed doc.update_status(DocStatus.PARSING) # 这里的 parse 在真实系统中应由 parser 完成 # 因为示例中 enqueue 时已经附带了 content,这里默认已经完成解析 doc.update_status(DocStatus.CHUNKING) if chunker and storage: chunks = chunker.split_document(doc) storage.upsert_chunks(doc.doc_id, chunks) # 标记最新的文档为已索引 storage.index_document(doc_id=doc.doc_id, source_name=doc.source_name) doc.update_status(DocStatus.INDEXED) print(f"[queue] 文档处理完成: {doc.source_name} (chunks: {len(chunks)})" if 'chunks' in locals() else f"[queue] 文档处理完成: {doc.source_name}") return True except Exception as e: doc.update_status(DocStatus.FAILED) print(f"[queue] 文档处理失败: {doc.source_name}, error={e}") return False def get_document(self, doc_id: str) -> Document | None: return self._documents.get(doc_id) def pending_count(self) -> int: return len(self._pending)队列实现虽然精简,但保留了状态机的核心逻辑。真实环境中,你可以在 process_next 中增加异常重试策略,或者在 enqueue 后立即将任务写入 Redis,由 Worker 进程异步消费。
4.5 实现存储与检索接口
存储层负责索引和检索 chunk。我们提供一个基于内存实现的版本,适合演示和学习。接口设计上预留了扩展点,方便后续接入向量数据库。
# 文件路径:docuqueue/storage.py from __future__ import annotations import math import re from collections import defaultdict from typing import Iterable from docuqueue.models import Chunk class Storage: """ 存储与检索接口。 当前实现基于简单的关键词匹配和 TF 分数排序,用于演示。 生产环境建议替换为向量检索或 Elasticsearch。 """ def __init__(self) -> None: self._chunks: dict[str, list[Chunk]] = defaultdict(list) self._documents: dict[str, dict] = {} def upsert_chunks(self, doc_id: str, chunks: list[Chunk]) -> None: self._chunks[doc_id] = chunks print(f"[storage] 索引 {len(chunks)} 个 chunks (doc_id={doc_id[:8]})") def index_document(self, doc_id: str, source_name: str) -> None: self._documents[doc_id] = {"doc_id": doc_id, "source_name": source_name} def search(self, query: str, top_k: int = 5) -> list[Chunk]: """ 基于关键词的词频评分进行简单检索。 """ query_terms = self._tokenize(query) if not query_terms: return [] scored_chunks: list[tuple[float, Chunk]] = [] for chunks in self._chunks.values(): for chunk in chunks: score = self._compute_score(chunk.content, query_terms) if score > 0: scored_chunks.append((score, chunk)) scored_chunks.sort(key=lambda x: x[0], reverse=True) return [chunk for _, chunk in scored_chunks[:top_k]] def list_documents(self) -> list[dict]: return list(self._documents.values()) def _compute_score(self, text: str, query_terms: list[str]) -> float: terms = self._tokenize(text) if not terms: return 0.0 term_count = len(terms) score = 0.0 for qt in query_terms: cf = terms.count(qt) if cf: # 简单的 TF * IDF 近似,IDF 部分这里简化为固定权重 score += cf / term_count return score def _tokenize(self, text: str) -> list[str]: """极简分词器:按非字母数字字符切分并统一小写。""" return [t.lower() for t in re.findall(r"\w+", text)] def get_chunks_by_doc(self, doc_id: str) -> list[Chunk]: return self._chunks.get(doc_id, [])这个检索实现的核心指标是“包含关键词的密度”,它虽然不如向量检索智能,但对于演示文档层流程绰绰有余。真实项目中,你可以把 search 方法替换为向量相似度计算,接口签名可以保持不变。
4.6 组装 DocuQueue 门面 API
门面类对上层隐藏内部细节,Agent 只需要调用一个简单的接口。
# 文件路径:docuqueue/docuqueue.py from __future__ import annotations from pathlib import Path from docuqueue.chunker import RecursiveCharacterChunker from docuqueue.models import Chunk, Document from docuqueue.parser import DocumentParser from docuqueue.queue import DocumentQueue from docuqueue.storage import Storage class DocuQueue: """ DocuQueue 门面类。 统一文档接入、解析、切分、索引和检索流程。 Agent 通过本类与文档层交互。 """ def __init__( self, chunk_size: int = 1000, chunk_overlap: int = 100, ) -> None: self.parser = DocumentParser() self.chunker = RecursiveCharacterChunker( chunk_size=chunk_size, chunk_overlap=chunk_overlap, ) self.storage = Storage() self.queue = DocumentQueue() self._init_queue_handlers() def _init_queue_handlers(self) -> None: # 可以注册额外的状态处理器,例如通知外部系统 pass def add_file(self, file_path: str | Path) -> str: """添加一个本地文件到文档队列,返回 doc_id。""" doc = self.parser.parse(file_path) self.queue.enqueue(doc) # 实际项目中,这里可以异步执行 self.queue.process_next( chunker=self.chunker, storage=self.storage, ) return doc.doc_id def add_text(self, content: str, source_name: str = "text", metadata: dict | None = None) -> str: """直接添加纯文本内容。""" doc = Document( source_name=source_name, content=content, metadata=metadata or {}, ) self.queue.enqueue(doc) self.queue.process_next( chunker=self.chunker, storage=self.storage, ) return doc.doc_id def search(self, query: str, top_k: int = 5) -> list[Chunk]: """在已索引的文档中检索相关内容。""" return self.storage.search(query, top_k=top_k) def get_document(self, doc_id: str) -> Document | None: return self.queue.get_document(doc_id) def list_documents(self) -> list[dict]: return self.storage.list_documents() def pending_count(self) -> int: return self.queue.pending_count()到这里,一个最小可用的文档层已经成型。它具备以下几个能力:
- 接受文件和纯文本输入
- 自动完成解析与切分
- 维护文档生命周期状态
- 支持关键词检索
- 返回与 RAG 兼容的 Chunk 结构
5. 将文档层接入 AI Agent 工作流
5.1 一个完整的 RAG 接入示例
下面的示例演示如何通过 OpenAI 兼容的 API 调用大模型,将 DocuQueue 检索到的文档片段作为上下文,生成回答。
# 文件路径:docuqueue/examples/demo_agent.py import os # 假设你的项目根目录在 docuqueue 上层 from docuqueue import DocuQueue def build_context(chunks) -> str: """把多个检索结果拼接成上下文文本。""" parts = [] for i, chunk in enumerate(chunks): header = f"[片段 {i + 1}] 来源: {chunk.metadata.get('source_name', 'unknown')}" parts.append(f"{header}\n{chunk.content}") return "\n\n".join(parts) def ask_agent(docuqueue: DocuQueue, question: str) -> str: """从文档层检索并向大模型发起对话请求。""" # 1. 检索相关文档片段 chunks = docuqueue.search(question, top_k=4) if not chunks: return "未在文档库中找到相关内容,请尝试其他问题。" # 2. 构建 Prompt context = build_context(chunks) prompt = f"""请根据以下文档内容回答问题。 文档内容: {context} 问题: {question} 请用简洁、准确的语言回答。如果文档内容不足以回答问题,请明确说明。""" # 3. 调用大模型接口(以 OpenAI 兼容接口为例) try: from openai import OpenAI except ImportError: return "未安装 openai 库,请先执行 pip install openai" client = OpenAI(api_key=os.getenv("OPENAI_API_KEY")) resp = client.chat.completions.create( model=os.getenv("OPENAI_MODEL", "gpt-4o-mini"), messages=[ {"role": "system", "content": "你是一个文档问答助手,善于利用给定的文档片段回答用户问题。"}, {"role": "user", "content": prompt}, ], temperature=0.3, ) return resp.choices[0].message.content def main() -> None: dq = DocuQueue(chunk_size=800, chunk_overlap=80) # 添加本地文档 dq.add_file("README.md") # 添加纯文本知识 dq.add_text( content="DocuQueue 是一个为 AI Agent 设计的文档层," "它统一管理文档的接入、解析、切分、索引和检索流程。" "通过队列机制实现异步处理和故障隔离," "让 Agent 可以专注于业务逻辑而不是文档细节。", source_name="intro.txt", ) print("文档索引完成。当前文档列表:") for doc in dq.list_documents(): print(f" - {doc['source_name']}") question = "DocuQueue 解决了什么问题?" print(f"\n问题: {question}") answer = ask_agent(dq, question) print(f"回答: {answer}") if __name__ == "__main__": main()运行示例前,请确认:
- 已设置 OPENAI_API_KEY 环境变量(或你使用的模型服务 API Key)。
- 已安装 openai 库。
- 当前目录存在 README.md 文件(如果没有,可以改为任意文本文件)。
5.2 与 Agent 框架的协作方式
在实际项目中,DocuQueue 通常作为 Agent 的“工具”存在。以 ReAct 或 Function Calling 模式为例,Agent 会接收到用户问题,然后决定调用哪个工具。文档问答场景下,Agent 会调用 search 工具,检索结果返回后,再连同问题一起交给模型生成答案。
在这种模式下,DocuQueue 并不关心 Agent 是用 LangChain、LlamaIndex 还是自研框架。真正重要的是它对外暴露了一套稳定接口:
- add_text / add_file 负责增量更新知识
- search 负责召回相关内容
- list_documents 便于 Agent 判断是否有足够的知识覆盖
6. 常见问题与排查思路
6.1 高频问题排查表
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 文档处理一直停在 pending | 队列没有消费者,或 process_next 未被调用 | 检查 Worker 进程是否在运行,确认队列消费逻辑被正确执行 |
| 检索结果不相关 | 切分粒度不合适,或检索算法过于简单 | 调小 chunk_size 尝试更细粒度;或升级为向量检索 |
| Token 超限 | 检索到的 chunk 数量过多 | 减小 top_k,或调整 chunk_size |
| 文档内容乱码 | 文件编码不是 UTF-8 | 在 parser 中指定编码,或使用 chardet 检测 |
| Agent 回答“未找到相关内容” | 检索没命中,或文档尚未完成索引 | 先调用 list_documents 确认索引状态;再检查提问关键词是否过于抽象 |
6.2 切分参数如何调整
chunk_size 和 chunk_overlap 是最常调优的两个参数。它们需要针对不同的文档类型做调整:
- 代码类文档:chunk_size 可以调大(1200-2000),因为代码块需要保持完整。
- 新闻或说明文:建议 chunk_size 在 500-1000,语义密度高,切小一点更精准。
- 法律或政策文件:每个条目可能很长,需要按章节切分,而不是固定窗口。
chunk_overlap 的作用是保证相邻 chunk 之间有信息的重叠区。如果你发现在回答中提到“上下文不连贯”,可以适当增加 overlap。
6.3 队列消费异常处理
真实环境中,consumer 进程可能因为各种原因崩溃。建议做三件事:
- 消费成功后显式提交偏移或删除任务,避免重复处理。
- 消费失败时记录 task_id 和错误堆栈,并写入死信队列。
- 为队列任务增加超时时间,防止某个文档卡住整个队列。
7. 最佳实践与工程建议
7.1 文档处理必须异步化
文档解析和切分属于 IO 密集型 + CPU 密集型操作。如果用户上传文档后需要等待整个流程结束,体验会非常差。建议把 add_file 和 add_text 设计为“入队即返回”,后台 worker 异步完成后续处理。Agent 可以通过查询接口轮询或接收 Webhook 通知。
7.2 元数据是检索质量的隐藏杠杆
很多人只关注文本切分,却忽略了元数据的重要性。丰富的元数据可以让检索系统在召回前先做过滤。例如:
- 为每个 chunk 标记文档的部门、时间、版本
- 在检索时支持 metadata filter
- 在回答时标注 chunk 来源,帮助用户判断可信度
7.3 先跑通再优化
不要一开始就追求复杂的向量检索、混合检索、重排模型。先用最简单的关键词检索跑通整个流程,确认业务逻辑是成立的,再逐步升级检索算法。每一步优化都要有评测数据支撑,否则很难判断是算法问题还是数据问题。
7.4 安全性:权限过滤不可跳过
如果你的文档库中包含私有或有权限限制的内容,检索必须做权限过滤。核心原则是:在检索阶段就把用户无权查看的内容排除掉。这意味着 chunk 的 metadata 中需要包含可见性标签,例如 team_id 或 role。生产环境中,文档权限是绝对不能省略的环节。
7.5 版本管理与更新策略
文档更新后,旧的 chunk 应该在索引中被替换。我们实现的 Storage 中,upsert_chunks 会根据 doc_id 直接覆盖旧 chunk,这是最简单可靠的做法。更复杂的场景下,你可能需要维护多个版本的文档快照,Agent 可以根据日期或版本号选择检索范围。
8. 总结与下一步学习路线
通过这篇教程,我们从零构建了一个面向 AI Agent 的 Document Layer,覆盖了几个关键点:统一文档模型、解析器注册机制、递归文本切分、文档队列状态管理、存储检索抽象,以及 Agent 接入示例。你可以在本地运行这段代码,体验从文档入队到问答响应的完整闭环。
如果想把这条技术路线继续深入,可以沿着下面几个方向往下走:
- 把 Storage 升级为向量数据库(如 Chroma、Milvus 或 qdrant),替换 search 方法的核心实现。
- 用 Redis Streams 替换内置 deque,实现多 Worker 并行消费。
- 引入重排模型,在召回后对 top_k 结果做精细排序。
- 为不同文件类型注册真实解析器,比如 PDF 解析、OCR 识别、表格还原。
DocuQueue 这类“文档层”真正解决的是 Agent 与数据之间的组织问题。它的核心价值不是某一种检索算法,而是让文档的接入、处理和读取变成一套可维护、可扩展的基础设施。如果你正在构建 Agent 应用,文档层绝对值得优先设计好,它会为你后续的迭代省下大量时间。