Pathway 实时 RAG 与 AG2 多智能体编排实战:让文档检索永远保持最新
本文基于仓库中开箱即用的示例 examples/projects/ag2-multiagent-rag/README.md,讲解如何把 Pathway 的实时文档索引能力(VectorStoreServer)作为 AG2(前 AutoGen)多智能体对话的检索后端:智能体通过工具调用查询最新知识库,从而在文档持续更新的场景下获得带引用的实时答案。读完本文,你将掌握「流式文档接入 + REST 向量检索服务 + 多智能体工具调用」的完整落地链路,并能直接运行该示例进行验证。
背景与核心思路
这个示例要解决的是一个在 RAG 应用中非常普遍的问题:传统做法把文档切分、向量化后写入静态索引,一旦源文档更新就必须手动重建索引,而知识库里的回答也随之过期。
示例的思路是把两条链路拼起来:
- Pathway 实时数据框架:用流式引擎持续读取、切分、向量化并增量索引文档,通过 VectorStoreServer 暴露 REST API;
- AG2:编排 Researcher(研究员)与 Analyst(分析师)等多个 AI 智能体,让智能体把检索后端当作工具(
search_documents)在对话过程中按需调用。
关键收益正如 README 强调的:Pathway 会在文档发生变化的瞬间自动重新索引,因此 AG2 智能体查询到的永远是知识库的最新版本,全程无需手动重建索引。这一组合尤其适合"文档频繁变动、智能体需要最新信息"的场景,例如实时更新的知识库、持续更新的报告或流式数据管道。
注:示例中使用的是 Pathway 与 AG2 两个框架的公共能力。AG2 是通用多智能体对话框架;Pathway 在本仓库中的核心实现位于 python/pathway/xpacks/llm/(LLM 与向量检索扩展)与 src/(Rust 流式引擎)。
系统架构与运行流程
示例 README 给出了清晰的端到端数据流:
Documents (live folder) --> Pathway Live Data Framework VectorStoreServer (real-time indexing)
|
REST API /v1/retrieve
|
User Query --> AG2 UserProxy --> GroupChat [Researcher + Analyst]
|
search_documents tool --> HTTP POST --> Pathway
|
Grounded, real-time answers with citations
对照示例脚本 main.py,整个流程可拆解为五个阶段:
pw.io.fs.read以**流式(streaming)**方式持续监听./data/目录(含子目录);- 新文档被切分为 token 级块,并经 OpenAIEmbedder 向量化,进入 VectorStoreServer 维护的实时索引;
- 后台线程启动 HTTP 服务,暴露
/v1/retrieve(相似度检索)等端点; - AG2 的 Researcher 智能体收到用户问题后,调用
search_documents工具向 Pathway 发起 HTTP POST 检索; - 检索结果连同来源信息回到群聊,Analyst 智能体汇总成带引用的结构化答案,最后以
TERMINATE结束对话。
环境准备
根据 README 与 requirements.txt,需要:
- Python >= 3.10;
- 一个 OpenAI API key(Pathway 端用于 embedding,AG2 端用于生成对话);
- 安装依赖:
pip install -U pathway "ag2[openai]>=0.11.4,<1.0" requests python-dotenv
依赖清单共四项,各自职责如下:
| 依赖 | 在本示例中的作用 |
|---|---|
pathway |
提供流式读取、OpenAIEmbedder、TokenCountSplitter 与 VectorStoreServer |
ag2[openai] |
多智能体框架及其 OpenAI 后端(版本约束 >=0.11.4,<1.0,同时兼容 autogen 命名空间导入) |
requests |
智能体侧向 Pathway /v1/retrieve 发起 HTTP 检索 |
python-dotenv |
从 .env 文件加载 OPENAI_API_KEY |
一个值得注意的细节:脚本中是以 from autogen import (...) 的方式导入 AG2 组件的(见 main.py),而安装包名为 ag2[openai]——这正是 AG2 为兼容前身 AutoGen API 而保留的命名空间,实际来源仍是所安装的 ag2 发行版。
配置:密钥、数据目录与许可证
密钥配置
在示例根目录创建 .env 文件并填入密钥即可:
OPENAI_API_KEY=your_openai_api_key_here
该变量会被两处消费:main.py 在启动时校验其是否存在(否则报错退出),随后既用于构造 AG2 的 LLMConfig,也会被 OpenAIEmbedder 自动读取(OpenAIEmbedder 的 docstring 明确说明 API key 可通过构造参数或 OPENAI_API_KEY 环境变量提供,见 embedders.py)。
数据目录
把 TXT、MD、PDF 等文档放进 ./data/,仓库已附带一份测试用样例文档 data/sample.md,其中分别介绍了 Pathway 与 AG2 的核心特性,恰好用于验证最终回答是否带引用。
脚本启动时会做两类目录检查:
- 若
./data/不存在则自动创建并提示放入文档后重跑; - 若目录为空(忽略隐藏文件),同样提示加入文档。
许可证设置
示例在启动时执行:
pw.set_license_key("demo-license-key-with-telemetry")
README 注明:如需使用 Pathway 的 Scale 版高级特性可申请免费许可证密钥填入此处;若使用 Community 版则注释掉该行即可。这属于运行环境相关配置,请根据你的实际授权方式决定是否保留。
运行示例
依次执行:
# 1. 放入文档(已附带 sample.md 时可直接运行)
python main.py
脚本会完成三件事:
- 在后台线程中启动 VectorStoreServer,实时索引
./data/中的文档; - 启动 AG2 多智能体(Researcher + Analyst + UserProxy),它们把 Pathway 索引作为工具使用;
- 打印带来源引用的多智能体对话结果。
由于示例同时运行 HTTP 服务与 AG2 对话,main.py 使用 threading.Thread(..., daemon=True) 启动服务端线程,避免相互阻塞。真正执行引擎的是 VectorStoreServer.run_server(见下文"服务与端点")。
就绪探测
服务是异步启动的,索引和 HTTP 监听都需要时间。脚本采用轮询就绪策略(main.py):最多 60 次、每 2 秒向 /v1/statistics 发起 POST,直到收到 200 响应并打印已索引文件数 file_count;60 次后仍未就绪则仅打印警告继续。这保证了 AG2 智能体开始对话时索引通常已可用。
运行中增删文档
README 特别强调:服务运行期间随时向 ./data/ 添加文档,Pathway 会自动完成重新索引,无需重启服务或手动触发。这正是实时 RAG 与"静态索引 + 手动重建"方案的本质区别。
深入原理:从配置到源码
1. 文档接入:流式读取本地目录
documents = pw.io.fs.read(
DATA_DIR,
format="binary",
mode="streaming",
with_metadata=True,
)
要点:
mode="streaming"使连接器持续监听目录变化(新增/修改/删除都会被感知),而非一次性读入后退出;format="binary"以原始字节交给后续解析器处理,因此同一份表格可容纳 TXT/MD/PDF 等不同格式;with_metadata=True保留文件路径等元信息,最终会在检索结果中回传,成为回答里"来源引用"的依据(如 sample.md 的path元数据)。
2. 向量化与切分组件
embedder = OpenAIEmbedder(model="text-embedding-3-small")
splitter = TokenCountSplitter(max_tokens=400)
OpenAIEmbedder是 xpacks.llm 对 OpenAI Embedding 服务的封装(embedders.py),构造函数支持capacity(并发上限)、retry_strategy(默认指数退避)、cache_strategy、batch_size、timeout等参数;本示例仅指定model="text-embedding-3-small"。TokenCountSplitter位于 splitters.py,按 token 数切分文档,max_tokens控制单块最大 token 数,其默认值为 500,示例调小为 400 以获得更细粒度、更利于定位的检索块;它还会尝试在合理位置断句,避免把语义割裂(参见该类的 docstring)。
3. VectorStoreServer:实时索引即服务
server = VectorStoreServer(
documents,
embedder=embedder,
splitter=splitter,
)
构造函数(vector_store.py)除文档表外还接受 parser、splitter、doc_post_processors 等参数,内部基于默认 KNN 索引构建检索器。
随后调用:
server.run_server(
host=PATHWAY_HOST, # "127.0.0.1"
port=PATHWAY_PORT, # 8765
threaded=False,
with_cache=True,
)
run_server(vector_store.py)负责把文档处理管线跑起来并绑定 HTTP 监听。关键参数:
| 参数 | 含义 |
|---|---|
host / port |
HTTP 监听地址与端口(本示例为 127.0.0.1:8765) |
threaded |
是否在线程中运行;True 返回线程对象,False 阻塞执行 |
with_cache |
是否对相同内容的 embedding 请求做缓存,默认开启 |
cache_backend |
缓存后端,默认本地磁盘目录 ./Cache(pw.persistence.Backend.filesystem),可通过 persistence API 覆盖 |
从源码看,服务实际注册了三个 REST 端点(vector_store.py):
/v1/retrieve:对查询做相似度检索,返回指定数量的文档;/v1/statistics:返回索引器当前统计(文档数与最近添加时间),示例用它做就绪探测;/v1/inputs:返回当前已索引文档的元数据列表。
这三个端点均通过 pw.io.http.rest_connector 接入实时计算图,因此任何时刻的查询结果都是索引当前状态的反映。当 with_cache=True 时,run_server 会以 PersistenceMode.UDF_CACHING 配置持久化(vector_store.py),避免重复内容反复产生 embedding 调用。
4. HTTP 检索:智能体视角的"搜索引擎"
query_pathway_server(main.py)演示了客户端侧完整调用:
url = f"http://{PATHWAY_HOST}:{PATHWAY_PORT}/v1/retrieve"
payload = {"query": query, "k": k}
response = requests.post(url, json=payload, timeout=30)
results = response.json()
返回体是 JSON 数组,每项含:
text:命中的文档块内容;metadata:来源元数据,其中path字段即源文件路径;dist:与查询的相似度距离,按 dist 升序排列(越小越相似)。
脚本将其格式化为 [i] Source: 路径 (distance: 值) + 正文 的多段文本返回给智能体,并针对服务未就绪(ConnectionError)与其它异常返回友好提示,避免工具异常中断群聊。
5. 多智能体编排与工具注册
AG2 侧由三类角色组成群聊(main.py):
- Researcher(
AssistantAgent):接到问题后调用search_documents检索知识库,若首次结果不足会换关键词重试,并给出带来源的发现; - Analyst(
AssistantAgent):基于 Researcher 的发现综合成条理清晰的答案,始终引用源文档,信息不足时反向要求 Researcher 换词检索,完成时输出TERMINATE; - UserProxy(
UserProxyAgent):human_input_mode="NEVER"全程无人介入,max_consecutive_auto_reply=10限制连续自动回复次数,code_execution_config=False禁止执行代码,并通过is_termination_msg识别含TERMINATE的消息以结束对话。
模型配置统一由 LLMConfig 提供,示例使用 gpt-4o-mini(main.py),可按你的账号权限替换。
Pathway 检索被注册为 Researcher 可调用的工具,是"多智能体 + 实时 RAG"结合的核心桥接代码:
@user_proxy.register_for_execution()
@researcher.register_for_llm(
description=(
"Search the document knowledge base powered by the Pathway Live Data Framework. "
"Returns relevant document chunks with source citations. "
"The knowledge base is continuously updated in real-time. "
"Use specific, targeted search queries for best results."
)
)
def search_documents(query: str, top_k: int = 5) -> str:
return query_pathway_server(query, k=top_k)
说明:
register_for_llm让 Researcher 的模型能看到并决定何时调用该工具,description是关键——它告诉模型"知识库实时更新、请用精准查询词";register_for_execution把同一函数交给 UserProxy 在本地实际执行,实现"LLM 决策、代理执行"的安全分离;- 默认
top_k=5,与检索端k参数对应。
群聊由 GroupChat 承载(agents 列表与消息历史),max_round=12 限制总轮次防止失控;GroupChatManager 负责调度发言顺序。最终入口 user_proxy.run(manager, message=query).process() 启动整场对话,query 要求智能体对文档关键主题做带引用的综合总结。
运行结果与验证建议
正常运行时终端会依次出现:VectorStoreServer 启动日志 → Pathway Live Data Framework server is ready! Indexed files: N(就绪探测成功)→ 默认查询文本 → 多智能体往返对话 → 以 TERMINATE 收尾的"Conversation complete."。
想要直观验证"实时索引"价值,可以这样做:
- 首次运行时对 data/sample.md 提问,观察回答是否引用该文件路径;
- 保持服务运行,向
./data/追加一段新文档,再次用相关问题提问——无需重启即可命中新内容; - 若需单独验证检索服务,可直接用
curl或任意 HTTP 客户端 POST 到http://127.0.0.1:8765/v1/retrieve,携带{"query": "...", "k": 5}查看原始返回结构(含text/metadata/dist字段)。
结语
本示例演示了一条可复用的模式:把 Pathway 的实时文档索引当作一组语义检索 REST 端点,供 AG2 等多智能体框架以工具方式在对话中按需调用。由于索引随文档流持续增量维护,智能体始终站在"最新知识"之上作答,且每个答案都携带可追溯的文档来源,天然适合实时知识库、滚动更新的报告以及流式数据驱动的问答场景。
若想进一步改造该示例,从仓库源码看有明确扩展点:将 pw.io.fs.read 替换为其它连接器(Kafka、S3、GDrive、PostgreSQL 等数据源均为 Pathway 原生支持);调整 TokenCountSplitter 的 max_tokens 或替换 OpenAIEmbedder 模型以适配不同语料粒度与成本;扩展 search_documents 的 description 或新增智能体角色以适配更复杂的研究流程。仓库 examples/projects/ 目录下还有更多现成的 AI 管道项目可供对照参考。
环境提示:本文涉及的命令与配置以当前仓库中的示例代码(main.py、requirements.txt)为准。运行前请确认 Python 版本 >= 3.10,并在
.env中正确配置OPENAI_API_KEY;若仅使用 Pathway Community 版,可注释掉pw.set_license_key(...)一行。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python07
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00