news 2026/9/8 23:57:44

使用 Pathway 构建可实时“遗忘”数据的 RAG 聊天机器人:Conf42 示例源码级解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 Pathway 构建可实时“遗忘”数据的 RAG 聊天机器人:Conf42 示例源码级解析

使用 Pathway 构建可实时“遗忘”数据的 RAG 聊天机器人:Conf42 示例源码级解析

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

本文基于仓库 examples/projects/conf42(对应文档 docs/2.developers/7.templates/ETL/_readmes/conf42.md)展开,讲解如何用 Pathway 从零搭建一个基于 RAG 的 LLM 问答机器人。它的核心亮点不是“回答问题”,而是让知识库与磁盘上的文档目录保持实时同步:当你新增或删除一份文档,索引与回答会立刻随之更新,机器人因此能过滤掉过时或错误的“假新闻”。读完本文,你将掌握其完整运行方式、main.py的每一条调用链,以及背后VectorStoreServerDocumentStore、解析器、切分器、Embedder 等 Pathway LLM xpacks 组件的真实实现原理。

这个 Demo 解决什么问题

传统 RAG 应用通常先离线建立索引,再反复查询同一个静态向量库。一旦底稿文档被修订、下线或撤回,索引若不重建,模型仍会引用过期内容作答——这在金融公告、实时资讯、内部知识管理等场景中可能带来严重误导。

Conf42 示例的做法不同:文档目录是一个活的数据源。目录中任何文件的出现、消失或改动都会作为数据流事件进入 Pathway 的增量计算引擎,向量索引随之更新;下一次查询时,被删除文档的向量和文本便不再参与召回与作答。实现对应仓库中的 examples/projects/conf42/main.py,示例文档包含 Alphabet 2023 年 2 月发布的 10-K 年报与 2023 年 6 月能源趋势报告(见 examples/projects/conf42/documents)。

运行环境与三步启动

按文档说明,运行该应用只需要三步:

# 1. 安装 Pathway pip install pathway # 2. 在 .env 文件中填入你的 OpenAI API Key # (代码中通过 load_dotenv() 读取,模型调用还需联网) # 3. 启动应用 python main.py

需要说明三点适用前提:

  • 本示例的向量索引与文本生成都依赖 OpenAI 接口(Embedder 与 Chat 模型),因此必须提供有效的OPENAI_API_KEY;在 main.py 中分别使用了text-embedding-ada-002嵌入模型和gpt-3.5-turbo对话模型。
  • 如果使用 Pathway 的 Community 版本,请注释掉 main.py 中的pw.set_license_key("demo-license-key-with-telemetry")一行;该行仅为启用 Scale 版高级特性而存在。
  • 若在本地跑 RAG 类流水线,通常建议安装带 LLM xpacks 完整依赖的版本,即pip install pathway[all](参见 docs/2.developers/7.templates/20.run-a-template.md),否则UnstructuredParserVectorStoreServer等模块所需的解析与向量依赖可能缺失。

启动后,应用会把 documents 目录内的 PDF 全部解析、切分、向量化并建立索引,随后在0.0.0.0:8000上监听 HTTP 请求。

用 curl 发起一次问答

文档给出了一个典型的 POST 请求。由于 rest_connector 的默认路由是/,请求直接发往根路径:

curl --data '{ "user": "user", "query": "What is the revenue of Alphabet in 2022 in millions of dollars?" }' http://localhost:8000/ | jq

请求体会被填入一个自定义的查询 Schema——query(问题)与user(用户名,可用来区分来源)。main.py中 PWAIQuerySchema 定义了这两个字段:

class PWAIQuerySchema(pw.Schema): query: str user: str

Pipeline 的处理逻辑是:查询到达后,先到向量索引里检索最相似的 1 篇文档(k=1),把文档文本与元数据拼进提示词,再交给 LLM 生成答案,最后把结果通过 REST 连接器写回给请求方。因此对上面这个问题,机器人会召回 Alphabet 10-K 年报中的相关段落并给出 2022 年营收数据。

关键:删除文档后“立即遗忘”

该示例最值得演示的操作是删除文档。文档原文给的例子是:

rm ./documents/20230203_alphabet_10K.pdf

一旦该文件从目录中消失,你接下来的所有查询都不会再引用这份文档。原因在于数据流向是活的,而不是一次性的:pw.io.fs.read作为输入连接器持续监听目录变化,删除文件会生成一条“移除”记录,增量计算引擎随即从索引中撤销该文档对应的所有向量与文本;后续召回时它不再出现,LLM 在提示词里也看不到它的内容,自然无法再“编造”出源自它的答案。同理,向目录中放入新的 PDF 也会自动完成解析、切分与索引,让机器人立刻学到新知识。

这正是演讲标题 "Make your LLM app sane again" 的含义——让 LLM 应用在现实世界数据变化发生时同步更新自身认知,而不是永远停留在首次建索引时刻的快照上。

main.py 逐步拆解:整条实时 RAG 链路

为了弄清楚“实时遗忘”是怎么实现的,我们把 examples/projects/conf42/main.py 从头到尾过一遍,并对照源码给出依据。

1. 以二进制流读取文档目录

documents = pw.io.fs.read("./documents/", format="binary", with_metadata=True)

format="binary"意味着输出的是原始字节而非文本;with_metadata=True会让每一条数据携带文件路径等元信息(后文拼提示词时用到了metadata["path"])。fs 输入连接器会把目录内容建模为一张持续演化的 Pathway 表,这正是“增删文档即增删数据”的入口。

2. 定义嵌入模型、对话模型与文本切分器

embedder = embedders.OpenAIEmbedder(model="text-embedding-ada-002") chat = llms.OpenAIChat(model="gpt-3.5-turbo", temperature=0.05) text_splitter = TokenCountSplitter(max_tokens=400)
  • OpenAIEmbedder是 xpacks 对 OpenAI Embedding 接口的封装(见 python/pathway/xpacks/llm/embedders.py),把文档段落变成向量,同时内置异步执行器、重试与缓存机制;
  • OpenAIChat封装 Chat Completion 接口(见 python/pathway/xpacks/llm/llms.py),temperature=0.05表示几乎确定性的低随机度输出,适合事实问答;
  • TokenCountSplitter按 token 数切分长文本,max_tokens=400控制每个片段的大小。源码位于 python/pathway/xpacks/llm/splitters.py,其max_tokens默认值为 500,这里显式收紧到 400 以控制单次送入模型的上文长度。

3. 用 VectorStoreServer 组装向量库

vector_server = VectorStoreServer( documents, embedder=embedder, splitter=text_splitter, parser=UnstructuredParser(), )

VectorStoreServer的类注释明确说明:它“构建一个文档索引流水线,并提供最近邻查询服务”(见 python/pathway/xpacks/llm/vector_store.py)。它接收输入文档表,内部依次完成三类处理:

  • 解析(parse)UnstructuredParser依赖 unstructured.io 把 PDF、DOCX 等二进制内容转成文本元素(见 python/pathway/xpacks/llm/parsers.py);
  • 切分(split):由传入的TokenCountSplitter把长文档切成便于嵌入的小块;
  • 向量化与索引(embed & index):构造时以DefaultKnnFactory(embedder=embedder)生成近邻检索器(见 vector_store.py)。

由于输入表本身是实时变化的,从解析到索引的整条链路都会随文件增删增量更新——这是“遗忘”能力的根本保证。

4. 通过 REST 连接器接收查询

webserver = pw.io.http.PathwayWebserver(host="0.0.0.0", port=8000) queries, writer = pw.io.http.rest_connector( webserver=webserver, schema=PWAIQuerySchema, autocommit_duration_ms=50, delete_completed_queries=True, )

rest_connector返回一对句柄:queries是承载每个 HTTP 请求的表,writer用于把计算结果写回对应响应。autocommit_duration_ms=50表示查询以 50ms 为批提交给流水线;delete_completed_queries=True表示处理完的请求会从表里清除,避免重复计算。若要深入了解这一 REST 接口设计,可参考 docs/2.developers/7.templates/40.rag-customization/10.REST-API.md。

5. 检索最相似的文档

results = queries + vector_server.retrieve_query( queries.select( query=pw.this.query, k=1, metadata_filter=pw.cast(str | None, None), filepath_globpattern=pw.cast(str | None, None), ) ).select( docs=pw.this.result, )

这里把查询表与retrieve_query的结果做 join,得到每个问题对应的召回文档列表。四个参数含义:

参数本示例取值作用
query用户问题作为检索的语义查询条件
k1返回最相似的文档数(此处只取 1 篇)
metadata_filterNone元数据过滤表达式,例如只检索某类路径的文档
filepath_globpatternNone按文件路径 glob 过滤可检索文档

retrieve_query的实现位于父类DocumentStore(见 python/pathway/xpacks/llm/document_store.py):它把查询交给_retriever.query_as_of_now(...),按number_of_matches=k做近邻搜索,最终以 JSON 列表形式返回{"text": ..., "metadata": ..., "dist": ...}metadata_filter在底层通过 jmespath 表达式对元数据求值过滤(见 vector_store.py)。

6. 把上下文拼成 RAG 提示词

@pw.udf def prep_rag_prompt(prompt: str, docs: list[pw.Json]) -> str: docs = docs.value docs = [{"text": doc["text"], "path": doc["metadata"]["path"]} for doc in docs] prompt_func = _unwrap_udf(prompts.prompt_short_qa) return prompt_func(prompt, docs)

这里把召回结果规整成text+path(来源文件路径)的结构,再调用prompts.prompt_short_qa生成提示词。prompt_short_qa位于 python/pathway/xpacks/llm/prompts.py,其模板指令值得注意:

  • “仅依据所给来源作答”(answer based solely on the provided sources),并要求回答简洁、准确;
  • 日期类问题需按严格格式输出、Yes/No 问题只答 Yes/No;
  • 如果问题无法从文档中推断,就输出No information found.

正是这最后一条规则,加上“删除文档后检索不到该文档”的事实,让机器人对已删除的 Alphabet 10-K 只会给出“找不到信息”,而非继续沿用旧记忆作答。_unwrap_udf只是把可能被装饰成pw.UDF的函数还原为原生可调用对象的小工具。

7. 交给 LLM 生成答案并写回

results += results.select(rag_prompt=prep_rag_prompt(pw.this.query, pw.this.docs)) results += results.select( result=chat( llms.prompt_chat_single_qa(pw.this.rag_prompt), ) ) writer(results) pw.run()

llms.prompt_chat_single_qa(见 python/pathway/xpacks/llm/llms.py)把单条问题字符串转换成 Chat 接口要求的[{role, content}]消息格式,再由OpenAIChat完成生成。最终writer(results)把答案列写回 REST 响应,pw.run()启动整个流式计算引擎并持续运行。

验证“实时遗忘”的实验路径

如果想在本地亲手验证行为变化,可按以下顺序操作:

  1. 启动服务后先查询 Alphabet 2022 年营收,能获得带数据来源的准确回答;
  2. 执行rm ./documents/20230203_alphabet_10K.pdf删除年报;
  3. 再次发送相同查询,观察返回结果变为“找不到相关信息”(No information found.)——旧文档不再被召回,也不再进入提示词;
  4. 反向实验:把删除的文件复制回目录(或放入新的 PDF),很快又可以检索到相关内容,说明新增文档也会被实时纳入索引。

每一步操作后通常需要等待数秒让引擎完成增量更新;如果观察不到变化,可先确认.env中 API Key 有效、pw.run()日志无报错,以及删除操作确实作用在./documents/目录内。

从示例到你的场景:可调整的扩展点

这个示例的骨架可以低成本复用到其他实时知识场景:

  • 检索条数:把retrieve_query中的k从 1 调大,让每次问答参考多篇文档,答案会更全面;可参考同目录下的 docs/2.developers/7.templates/ETL/_readmes/question-answering-rag.md 对比不同 RAG 模板的取舍;
  • 文档来源VectorStoreServer接受任何 Pathway 表作为输入(构造签名是*docs: pw.Table),因此把 fs 目录读取替换为 Kafka、S3、数据库或对象存储连接器,即可对实时事件流、云端文件做同样的“文档即数据流”RAG;
  • 元数据过滤:传入metadata_filter(jmespath 表达式)或filepath_globpattern,可在召回阶段按来源路径、标签等条件圈定可检索范围;
  • 提示策略:替换prompts.prompt_short_qaprompt_qaprompt_rerank或自定义模板(全部定义在 python/pathway/xpacks/llm/prompts.py),即可改变回答风格、加入引用来源或强制格式约束。

小结

Conf42 示例演示了一个很容易被忽视却影响巨大的 RAG 事实:当知识源发生增删时,问答系统应当随之更新,而不是守着过期的静态索引。通过把“文档目录”建模为 Pathway 的实时输入,并把解析、切分、嵌入、近邻检索整条链路都构建在增量数据流之上,main.py 用约 80 行代码实现了“文件一删,模型即忘”的能力。理解这一模式后,你可以把同样的思想迁移到新闻真伪过滤、财报审计、法规更新追踪等任何需要“知识与现实同步”的 LLM 应用中去。

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/8 23:57:36

ComfyUI 工作流导入导出完整指南:3 步保存、分享与复现

ComfyUI 工作流导入导出完整指南:3 步保存、分享与复现 【免费下载链接】ComfyUI The most powerful and modular diffusion model GUI, api and backend with a graph/nodes interface. 项目地址: https://gitcode.com/GitHub_Trending/co/ComfyUI ComfyUI …

作者头像 李华
网站建设 2026/9/8 23:55:47

Python人脸识别系统实战:从OpenCV到face_recognition的完整指南

简介:一套基于Python的人脸识别系统完整工程,适合Python初学者、计算机视觉入门者以及需要快速搭建人脸识别Demo的开发者。资源覆盖从摄像头人脸采集、特征提取、数据库建库到实时识别比对的整套流程,并配有tkinter图形界面与运行说明文档&am…

作者头像 李华
网站建设 2026/9/8 23:51:47

业务分析师快速了解陌生项目:从认知清单到干系人地图的实操路径

接到通知去支援一个新项目,第一天坐进工位,面对一堆陌生的系统名称、业务术语、项目文档,开发催着要需求,业务方等着确认流程,直属领导问“你了解得怎么样了”。这种场景,做过BA的朋友多半都经历过。BA全称…

作者头像 李华
网站建设 2026/9/8 23:51:45

002 — Globex GmbH — Staff DevOps Engineer

002 — Globex GmbH — Staff DevOps Engineer 【免费下载链接】career-ops Open-source AI job search: scan job portals, evaluate listings into a structured A-H report with a global 1-5 score, tailor your CV, track applications — runs locally in your AI coding…

作者头像 李华
网站建设 2026/9/8 23:50:42

XL5301 dToF传感器深度解析:宽电压、低功耗、高稳定性实战指南

1. 项目概述:为什么XL5301一出来,我就立刻拆了三颗样片上电测试TOF传感器这个圈子其实很小,老玩家基本都用过XL5300——它在2020年前后是国产dToF方案里少有的能稳定做到2.5米10%反射率、功耗压到8mA10Hz的型号,被大量用在扫地机避…

作者头像 李华
网站建设 2026/9/8 23:50:01

ip2region 完整指南:3 步搞定离线 IP 定位到城市级

ip2region 完整指南:3 步搞定离线 IP 定位到城市级 【免费下载链接】ip2region Ip2region is an offline IP-to-Region localization library and IP data management framework with both IPv4 and IPv6 supports, 10-microsecond level query efficiency, xdb sea…

作者头像 李华