首页
/ Pathway 实时 RAG 与 AG2 多智能体编排实战:让文档检索永远保持最新

Pathway 实时 RAG 与 AG2 多智能体编排实战:让文档检索永远保持最新

2026-09-07 16:39:26作者:蔡丛锟

本文基于仓库中开箱即用的示例 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,整个流程可拆解为五个阶段:

  1. pw.io.fs.read 以**流式(streaming)**方式持续监听 ./data/ 目录(含子目录);
  2. 新文档被切分为 token 级块,并经 OpenAIEmbedder 向量化,进入 VectorStoreServer 维护的实时索引;
  3. 后台线程启动 HTTP 服务,暴露 /v1/retrieve(相似度检索)等端点;
  4. AG2 的 Researcher 智能体收到用户问题后,调用 search_documents 工具向 Pathway 发起 HTTP POST 检索;
  5. 检索结果连同来源信息回到群聊,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 提供流式读取、OpenAIEmbedderTokenCountSplitterVectorStoreServer
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

脚本会完成三件事:

  1. 后台线程中启动 VectorStoreServer,实时索引 ./data/ 中的文档;
  2. 启动 AG2 多智能体(Researcher + Analyst + UserProxy),它们把 Pathway 索引作为工具使用;
  3. 打印带来源引用的多智能体对话结果。

由于示例同时运行 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_strategybatch_sizetimeout 等参数;本示例仅指定 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)除文档表外还接受 parsersplitterdoc_post_processors 等参数,内部基于默认 KNN 索引构建检索器。

随后调用:

server.run_server(
    host=PATHWAY_HOST,      # "127.0.0.1"
    port=PATHWAY_PORT,      # 8765
    threaded=False,
    with_cache=True,
)

run_servervector_store.py)负责把文档处理管线跑起来并绑定 HTTP 监听。关键参数:

参数 含义
host / port HTTP 监听地址与端口(本示例为 127.0.0.1:8765
threaded 是否在线程中运行;True 返回线程对象,False 阻塞执行
with_cache 是否对相同内容的 embedding 请求做缓存,默认开启
cache_backend 缓存后端,默认本地磁盘目录 ./Cachepw.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_servermain.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(UserProxyAgenthuman_input_mode="NEVER" 全程无人介入,max_consecutive_auto_reply=10 限制连续自动回复次数,code_execution_config=False 禁止执行代码,并通过 is_termination_msg 识别含 TERMINATE 的消息以结束对话。

模型配置统一由 LLMConfig 提供,示例使用 gpt-4o-minimain.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."。

想要直观验证"实时索引"价值,可以这样做:

  1. 首次运行时对 data/sample.md 提问,观察回答是否引用该文件路径;
  2. 保持服务运行,向 ./data/ 追加一段新文档,再次用相关问题提问——无需重启即可命中新内容;
  3. 若需单独验证检索服务,可直接用 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 原生支持);调整 TokenCountSplittermax_tokens 或替换 OpenAIEmbedder 模型以适配不同语料粒度与成本;扩展 search_documentsdescription 或新增智能体角色以适配更复杂的研究流程。仓库 examples/projects/ 目录下还有更多现成的 AI 管道项目可供对照参考。

环境提示:本文涉及的命令与配置以当前仓库中的示例代码(main.pyrequirements.txt)为准。运行前请确认 Python 版本 >= 3.10,并在 .env 中正确配置 OPENAI_API_KEY;若仅使用 Pathway Community 版,可注释掉 pw.set_license_key(...) 一行。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.74 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.81 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
531
595
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
920
1.84 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.63 K
1.02 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
518
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
389