使用 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的每一条调用链,以及背后VectorStoreServer、DocumentStore、解析器、切分器、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),否则UnstructuredParser、VectorStoreServer等模块所需的解析与向量依赖可能缺失。
启动后,应用会把 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: strPipeline 的处理逻辑是:查询到达后,先到向量索引里检索最相似的 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 | 用户问题 | 作为检索的语义查询条件 |
k | 1 | 返回最相似的文档数(此处只取 1 篇) |
metadata_filter | None | 元数据过滤表达式,例如只检索某类路径的文档 |
filepath_globpattern | None | 按文件路径 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()启动整个流式计算引擎并持续运行。
验证“实时遗忘”的实验路径
如果想在本地亲手验证行为变化,可按以下顺序操作:
- 启动服务后先查询 Alphabet 2022 年营收,能获得带数据来源的准确回答;
- 执行
rm ./documents/20230203_alphabet_10K.pdf删除年报; - 再次发送相同查询,观察返回结果变为“找不到相关信息”(
No information found.)——旧文档不再被召回,也不再进入提示词; - 反向实验:把删除的文件复制回目录(或放入新的 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_qa为prompt_qa、prompt_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),仅供参考