首页
/ Pathway 实时索引 + AG2 多智能体 RAG:解析 ag2-multiagent-rag 示例知识文档及完整落地实现

Pathway 实时索引 + AG2 多智能体 RAG:解析 ag2-multiagent-rag 示例知识文档及完整落地实现

2026-09-07 10:39:48作者:田桥桑Industrious

导读:examples/projects/ag2-multiagent-rag/data/sample.md 是仓库中「AG2 多智能体 × Pathway Live Data Framework 实时 RAG」示例内置的样例知识文档,用于验证多智能体在实时更新的知识库上检索与作答的全流程。本文将以此为骨架,讲清这份文档描述的两大框架核心能力,并结合同目录 README.mdmain.pypathway.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.pystart_pathway_server()

documents = pw.io.fs.read(
    DATA_DIR,
    format="binary",
    mode="streaming",      # 流式监控目录,新文件/改动会被持续发现
    with_metadata=True,    # 保留文件路径等元数据,供检索结果引用
)

pw.io.fs.readmode="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.readpw.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_strategyapi_keybatch_sizetruncation_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.pyrun_serverpython/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 还接受 threadedwith_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 构造了一个三角色 GroupChatexamples/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(LLMConfigmodel="gpt-4o-mini"api_type="openai"),但 AG2 的 LLMConfig 抽象同样面向其他 provider;而"灵活的编排"体现在 GroupChat + GroupChatManager 之上:研究员先检索、分析师再归纳、知识不足时分析师可让研究员换关键词再搜,最终由终止条件收口。

四、检索与引用链:Pathway 返回什么、智能体怎么用

query_pathway_servermain.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。

  1. 安装依赖(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]

  2. 配置密钥:在工程根目录创建 .env 文件写入 OPENAI_API_KEY=your_openai_api_key_here;main.py 通过 load_dotenv() 读取,并在缺失时直接报错退出。

  3. 准备语料:把文档放入 ./data/(示例自带 sample.md)。main.py 会先检查目录存在且非空。

  4. 运行:

    python main.py
    

    脚本依次执行:以守护线程启动 VectorStoreServer → 每 2 秒轮询 /v1/statistics 直到返回 200(并打印 file_count)→ 创建 AG2 角色与工具 → 发起总结性问题并打印完整多智能体对话。

  5. 运行期间向 ./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.mdmain.pyvector_store.py 的源码,读者可以以此为模板,把数据源换成 Kafka/S3/PostgreSQL,把语料换成自身业务文档,快速搭建"文档频繁变动、智能体始终读到最新事实"的生产级实时 RAG 问答系统。

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