Pathway 实时索引 + AG2 多智能体 RAG:解析 ag2-multiagent-rag 示例知识文档及完整落地实现
导读:
examples/projects/ag2-multiagent-rag/data/sample.md是仓库中「AG2 多智能体 × Pathway Live Data Framework 实时 RAG」示例内置的样例知识文档,用于验证多智能体在实时更新的知识库上检索与作答的全流程。本文将以此为骨架,讲清这份文档描述的两大框架核心能力,并结合同目录 README.md、main.py 与pathway.xpacks.llm源码,说明文档增改后如何被自动重新索引、智能体如何通过/v1/retrieve拿到带引用的检索结果,让读者既能读懂样例数据,也能照此搭建自己的实时多智能体问答管线。
一、先认识这份文档:它不是教程,而是「被检索的知识库」
data/sample.md 文件(examples/projects/ag2-multiagent-rag/data/sample.md)本身只有 28 行,标题为「Pathway Live Data Framework + AG2 Sample Document」,其定位是示例工程的测试语料:当运行 python main.py 时,Pathway 会持续监控 ./data/ 目录,把其中的 TXT/MD/PDF 文档实时解析、切分、向量化并建索引;而这份 sample.md 正是 AG2 智能体发起的首个问题("What are the key topics and insights described in the documents?")所对应的答案来源。
因此,理解这份文档的正确姿势是把它当作「知识库内容」去剖析——它概括了本项目涉及的两大技术栈(Pathway 与 AG2)的核心卖点,而这些卖点恰好又是整个示例工程逐条兑现的功能清单。后续章节将逐条对照源码验证。
二、sample.md 中的 Pathway 特性:逐条对照源码验证
文档第二节「About the Pathway Live Data Framework」列出了关于 Pathway 的五条判断,下面逐一与仓库源码比对。
1. "Python ETL framework … with a high-performance Rust engine under the hood"
Pathway 对外提供 Python API,但核心执行引擎由 Rust 实现。这在工程目录结构上有直观印证:仓库顶层 src/ 目录下保存了 139 个 Rust 源文件,与提供 Python 绑定的 python/ 目录并存。文档提到的 multithreading / multiprocessing 能力,则由底层数据流运行时按需并行调度实现。对使用方而言,这一分层意味着可以"用 Python 的易用性写逻辑、用 Rust 的吞吐跑数据"。
2. "Real-time document indexing and re-indexing" 与 "automatic handling of document updates"
这是整个示例的灵魂:无需人工重新索引,文档一变,索引立刻跟着变。源码证据在 main.py 的 start_pathway_server():
documents = pw.io.fs.read(
DATA_DIR,
format="binary",
mode="streaming", # 流式监控目录,新文件/改动会被持续发现
with_metadata=True, # 保留文件路径等元数据,供检索结果引用
)
pw.io.fs.read 以 mode="streaming" 读取目录后,数据源被接入一个持续增量计算的引擎:文件新增、删除、修改都会作为数据流中的变更事件被处理,驱动下游切分与向量化步骤自动重算,从而实现文档的实时(re)indexing——这正是"自动处理文档更新、无需手动重建索引"的机制来源。
3. "Support for various data sources (Kafka, S3, GDrive, PostgreSQL, and more)"
目录监控只是 Pathway 连接器(connectors)中的一种。仓库中 integration_tests/ 下存在成体系的连接器集成测试,例如:
integration_tests/kafka/:Kafka 流接入测试;integration_tests/s3/与integration_tests/db_connectors/(含 PostgreSQL 场景);integration_tests/gdrive/、integration_tests/sharepoint/:网盘/协作文档类源。
sample.md 中"支持多种数据源"的表述,在工程层面即有这些可运行的接入与测试作为支撑。实际应用中,只需把 main.py 里的 pw.io.fs.read 替换为对应连接器(如 pw.io.kafka.read、pw.io.s3.read),后续的切分、嵌入、建索引、检索链路完全复用。
4. "Built-in LLM integration with the xpacks.llm module"
示例从 pathway.xpacks.llm 导入了三个核心组件,正是文档所讲的"内置 LLM 集成":
from pathway.xpacks.llm.embedders import OpenAIEmbedder
from pathway.xpacks.llm.splitters import TokenCountSplitter
from pathway.xpacks.llm.vector_store import VectorStoreServer
- 嵌入器 OpenAIEmbedder:OpenAI 向量服务的封装。示例使用
text-embedding-3-small;源码实现支持capacity(最大并发)、retry_strategy(默认指数退避)、cache_strategy、api_key、batch_size、truncation_keep_strategy等参数,且api_key既可显式传入也可从OPENAI_API_KEY环境变量读取。 - 切分器 TokenCountSplitter:按 token 数切分长文档。默认
min_tokens=50, max_tokens=500,采用 tiktoken 编码(默认cl100k_base),并"尽量在标点处断句"——源码中维护PUNCTUATION = [".", "?", "!", "\n"]列表,切出的块会回退到最近标点,避免把语义截断。示例代码将max_tokens收紧为 400,便于检索粒度更细。 - 索引服务 VectorStoreServer:见下节。
5. "VectorStoreServer for serving document embeddings via REST API"
示例通过 VectorStoreServer 把"文档表 + 嵌入器 + 切分器"组合成一条实时索引管线,再经 run_server() 对外提供 REST 服务。其底层(vector_store.py 中 run_server,python/pathway/xpacks/llm/vector_store.py#L64-L131)用 pw.io.http.rest_connector 挂载了多个内建路由,与文档及 main.py 直接相关的有:
| 路由 | 作用 | 示例中的用法 |
|---|---|---|
/v1/retrieve |
对实时维护的索引做相似度检索,返回 text + metadata + dist |
智能体的 search_documents 工具 POST 查询 |
/v1/statistics |
返回索引器当前统计(如已索引文件数 file_count) |
main.py 等待服务就绪的健康检查 |
/v1/inputs |
返回当前索引中的文档列表(元数据) | 调试知识库内容时使用 |
run_server 还接受 threaded 与 with_cache 参数:示例中以 threaded=False 把服务跑在独立后台线程,with_cache=True 使相同文本的嵌入请求可命中磁盘缓存(默认 ./Cache,可用 cache_backend 覆盖),避免重复调用 OpenAI 计费。
三、sample.md 中的 AG2 特性:对照 main.py 的编排实现
文档第三节介绍 AG2(原 AutoGen)的多智能体能力,main.py 是对这几条特性的直接编码实现。
1. "Multi-agent conversations with GroupChat"
main.py 构造了一个三角色 GroupChat(examples/projects/ag2-multiagent-rag/main.py#L234-L243):
group_chat = GroupChat(
agents=[user_proxy, researcher, analyst],
messages=[],
max_round=12, # 对话最多 12 轮,防止无限循环
)
manager = GroupChatManager(groupchat=group_chat, llm_config=llm_config)
researcher(研究员):拿到问题后调用search_documents工具检索知识库,负责"找证据";analyst(分析师):基于研究员带回的片段做综合归纳,要求始终引用来源文档,并在答案完整时以TERMINATE结束;user_proxy:人类代理端,human_input_mode="NEVER"、max_consecutive_auto_reply=10,并用is_termination_msg识别TERMINATE信号来结束会话。
2. "Tool registration for agents via decorator pattern"
search_documents 是 AG2 装饰器式工具注册的典型范例(main.py#L212-L231):
@user_proxy.register_for_execution() # 由 UserProxy 负责实际执行
@researcher.register_for_llm(description=...) # 向 researcher 的 LLM 暴露工具
def search_documents(query: str, top_k: int = 5) -> str:
return query_pathway_server(query, k=top_k)
外层装饰器把工具绑定给 LLM(并附上"知识库实时更新、请用具体查询词"的描述,帮助模型正确调用),内层装饰器把函数交给代理执行;请求从「研究员 LLM 决定调用」到「UserProxy 发起 HTTP POST」再到「结果以带引用文本回填对话」的闭环由此完成。
3. "Support for various LLM providers" 与 "flexible agent orchestration"
示例使用 OpenAI(LLMConfig 中 model="gpt-4o-mini"、api_type="openai"),但 AG2 的 LLMConfig 抽象同样面向其他 provider;而"灵活的编排"体现在 GroupChat + GroupChatManager 之上:研究员先检索、分析师再归纳、知识不足时分析师可让研究员换关键词再搜,最终由终止条件收口。
四、检索与引用链:Pathway 返回什么、智能体怎么用
query_pathway_server(main.py#L82-L119)向 http://127.0.0.1:8765/v1/retrieve POST {"query": ..., "k": top_k},并把响应拼装为带来源的检索片段。其注释明确给出了 /v1/retrieve 的返回结构:
[{"text": "...", "metadata": {...}, "dist": float}, ...]
text:命中的文档分块;metadata.path:来源文件路径,用于生成引用(Source: xxx.md);dist:查询与分块的距离,值越小越相似,返回结果按dist升序排列。
正是 metadata 里保留的 path,让 Analyser 能输出"带出处、可溯源"的答案;而因为 Pathway 侧的索引实时维护,同一文件被编辑后再次检索,返回的分块内容即为最新版本——这正是本文开头所述"多智能体永远查询到最新知识库"的关键。
五、从零运行:让 sample.md 真正被"读进去"
整体数据流(摘自 README.md)如下:
Documents (live folder) --> Pathway 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
运行前置条件:Python >= 3.10,且需要 OpenAI API key。
-
安装依赖(
requirements.txt与 main.py docstring 内容一致):pip install -U pathway "ag2[openai]>=0.11.4,<1.0" requests python-dotenv注意:示例源码仍以
from autogen import ...导入 AG2 组件(AG2 沿用了 AutoGen 的导入名),因此安装包名为ag2[openai]。 -
配置密钥:在工程根目录创建
.env文件写入OPENAI_API_KEY=your_openai_api_key_here;main.py 通过load_dotenv()读取,并在缺失时直接报错退出。 -
准备语料:把文档放入
./data/(示例自带 sample.md)。main.py 会先检查目录存在且非空。 -
运行:
python main.py脚本依次执行:以守护线程启动 VectorStoreServer → 每 2 秒轮询
/v1/statistics直到返回 200(并打印file_count)→ 创建 AG2 角色与工具 → 发起总结性问题并打印完整多智能体对话。 -
运行期间向
./data/追加或修改文档:Pathway 流式监控会自动发现变更并重算索引,无需重启或手动重建——这就是 sample.md 第 1、5 条特性的现场验证方式。
关于许可证:main.py 顶部调用了 pw.set_license_key("demo-license-key-with-telemetry"),其注释说明——若使用 Pathway Scale 高级特性可替换为免费申请到的 license key;若使用 Community 版则注释掉该行即可。
六、小结:这份样例文档教你什么
sample.md 篇幅虽短,却是一张"功能验收清单":它把 Pathway 的实时索引、多数据源、xpacks.llm 集成、VectorStoreServer REST 服务、免手动重索引,以及 AG2 的 GroupChat、装饰器式工具注册、多 LLM 与灵活编排,全部浓缩为可供多智能体检索的测试语料。配合 README.md、main.py 与 vector_store.py 的源码,读者可以以此为模板,把数据源换成 Kafka/S3/PostgreSQL,把语料换成自身业务文档,快速搭建"文档频繁变动、智能体始终读到最新事实"的生产级实时 RAG 问答系统。
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 StartedRust0627
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00