Pathway 实战:用 Pathway 构建可实时更新的检索增强生成(RAG)问答流水线
导读
本文围绕仓库中的 question-answering-rag 模板说明 及其配套的 main.py,完整讲解如何基于 Pathway 的 LLM xpack(pathway[xpack-llm])从零搭建一条"文档问答式"的检索增强生成(Retrieval-Augmented Generation, RAG)流水线。你会掌握从文档解析、切分、向量索引,到查询检索、上下文构建、LLM 生成答案、HTTP 接口暴露的完整链路,并理解 Pathway 使这条链路"默认实时"——当索引的文档发生变化时,索引与问答服务会自动增量更新,无需重启。
模板概述:一个自带 HTTP 服务的 RAG 问答服务
该示例位于 examples/projects/question-answering-rag 目录,通过"检索 + 生成"的组合提升回答的准确性与相关性:先召回与用户问题最相关的文档片段,再让大语言模型基于这些片段生成答案。这样既保证了生成的连贯性,也让回答内容有据可依、贴近最新文档。
模板由两部分组成:
- RAG 流水线本体(main.py):负责文档摄取、向量索引、查询检索与答案生成;
- 轻量 HTTP 服务器:把流水线包装成一个 Web 服务,用户可用
curl直接发送查询并取回回答,默认监听地址0.0.0.0:8011。
同一份模板说明在文档区
docs/2.developers/7.templates/ETL/_readmes/question-answering-rag.md也有镜像副本,二者内容一致;本仓库根目录下的路径以 examples/projects/question-answering-rag/main.py 为可直接运行的实际源码。
RAG 流水线的总体架构
模板 README 对 RAG 的结构给出了清晰的七段式拆解,这也是理解 main.py 中每个函数对应职责的纲领:
- 文档索引(Document Indexing):流程起始于一批文档,先对其内容进行分析、抽取关键信息,并组织成可检索的索引结构,以便后续高效检索。
- 用户查询(User Query):用户输入一个问题或信息请求,作为 RAG 流程的起点。
- 文档检索(Document Retrieval):检索系统拿着用户查询在索引中搜索最相关的信息片段,利用向量相似度等算法快速定位可能包含答案的文档。
- 上下文构建(Context Building):把检索到的文档整理成上下文,作为后续生成回复的依据。
- 提示词构造(Prompt Construction):把"用户查询 + 检索得到的上下文"拼装成一个提示词(prompt),作为生成式语言模型的输入。
- 答案生成(Answer Generation):大语言模型处理该提示词,生成连贯且准确的回复。
- 最终输出(Final Output):把生成的回答返回给用户。
将检索与生成相结合后,回答不仅准确,而且与上下文强相关——这正是 RAG 适合复杂信息检索问答场景的原因。在 Pathway 中,文档索引与检索的"核心仓库"由 DocumentStore 类承担,定义在 python/pathway/xpacks/llm/document_store.py 中,下面的源码拆解会逐步落到该类及其周边组件上。
环境准备与安装
前提条件
按模板要求,先安装带 LLM xpack 扩展的 Pathway 与 python-dotenv:
pip install pathway[xpack-llm] python-dotenv
其中:
pathway[xpack-llm]是带 AI 组件(DocumentStore、UnstructuredParser、各类 embedder 与 LLM Chat 封装等)的完整发行版;python-dotenv用于从项目根目录的.env文件中加载OPENAI_API_KEY。
获取示例代码
在本仓库中,示例代码就在 examples/projects/question-answering-rag 目录,包含两个文件:
README.md:模板使用说明;main.py:完整的 RAG 流水线实现。
进入示例目录:
cd examples/projects/question-answering-rag/
配置环境变量
在项目根目录创建 .env 文件并填入 OpenAI API Key:
OPENAI_API_KEY=your_openai_api_key_here
main.py 顶部通过 load_dotenv() 加载该文件(见 main.py)。API Key 随后被读取并传给两个组件:文档向量化用的 OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"]),以及生成回答用的 llms.OpenAIChat(..., api_key=os.environ["OPENAI_API_KEY"])。从源码结构看,OpenAIChat 也支持通过同名的 OPENAI_API_KEY 环境变量自动读取 Key,无需显式传参(见 llms.py 中 OpenAIChat 的构造说明)。
数据准备与 HTTP 服务配置
待索引文档放置位置
流水线的文档入口是 get_documents():
def get_documents() -> DocumentSchema:
documents = pw.io.fs.read("./data/", format="binary", with_metadata=True)
return documents
它用 pw.io.fs.read 以二进制格式读取本地 ./data/ 目录,并开启 with_metadata=True(把文件路径等元数据一并带入下游,便于检索结果回溯来源)。因此你需要把待问答的文档(PDF、DOCX、TXT、Markdown、HTML 等)放进该目录。
HTTP 服务器
模板自带一个轻量 HTTP 服务器用于收发查询,默认 host="0.0.0.0"、port=8011,可通过修改这一行调整:
webserver = pw.io.http.PathwayWebserver(host="0.0.0.0", port=8011)
与之配套的接入方式在 get_webserver():
webserver = pw.io.http.PathwayWebserver(host="0.0.0.0", port=8011)
queries, writer = pw.io.http.rest_connector(
webserver=webserver,
schema=QuerySchema,
autocommit_duration_ms=50,
delete_completed_queries=False,
)
rest_connector 把 HTTP POST 请求体按 QuerySchema 解析成一张 Pathway 实时表(表内每行是一次查询),并返回一个 writer,用于把最终回答写回 HTTP 响应。两个关键参数:
autocommit_duration_ms=50:收到的查询以约 50ms 的批处理窗口提交进入流水线,兼顾吞吐与低延迟;delete_completed_queries=False:处理完成的查询行不从表中删除,方便调试与观察。
传入请求体的 schema 定义在 main.py:
class QuerySchema(pw.Schema):
messages: str
也就是说,请求体只需一个字段 messages,即用户的原始问题文本。
源码拆解:从文档到索引的摄取链路
RAG 的第一步是把 ./data/ 中的原始文档转成可检索的向量索引,这条链路在 get_documentstore() 中组装完成,核心是 Pathway LLM xpack 的 DocumentStore:
def get_documentstore(documents: DocumentSchema) -> DocumentStore:
text_splitter = TokenCountSplitter(
min_tokens=100, max_tokens=500, encoding_name="cl100k_base"
)
embedder = OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"])
retriever_factory = BruteForceKnnFactory(embedder=embedder)
parser = UnstructuredParser(
chunking_mode="by_title",
chunking_kwargs={
"max_characters": 3000,
"new_after_n_chars": 2000,
},
)
store = DocumentStore(
docs=documents,
retriever_factory=retriever_factory,
parser=parser,
splitter=text_splitter,
)
return store
各组件职责如下:
1. 解析器:UnstructuredParser
解析工作委托给 unstructured 库。模板把 chunking_mode 设为 "by_title",即按章节/标题边界切块,同时保留段落语义:
"by_title":在 basic 策略基础上按文档标题切块,保留章节边界(更利于"按文档章节回答问题");- 其他可用模式还包括
"basic"(默认的通用切块)、"single"(整篇作为一个长文本)、"elements"(按 Unstructured 元素拆分),校验逻辑见 parsers.py 中UnstructuredParser的定义。
chunking_kwargs 中的 max_characters=3000 与 new_after_n_chars=2000 用于限制单块字符数与换块时机,防止切出过长的文本块。
2. 切分器:TokenCountSplitter
即便解析器已经按标题切块,模板仍用 TokenCountSplitter 按 token 数做二次精切:min_tokens=100、max_tokens=500、按 OpenAI 的 cl100k_base 编码(即 GPT-4 系 / text-embedding 系模型使用的词表)统计 token,保证送入向量模型的每个块在长度上稳定可控。该组件定义在 splitters.py。
3. 向量化器:OpenAIEmbedder
OpenAIEmbedder 把切分后的文本块编码成稠密向量。模板只传了 api_key,模型参数使用默认值 "text-embedding-3-small"(默认值定义在 embedders.py)。text-embedding-3-small 的向量维度为 1536、最大 token 限制为 8191,相关模型与维度的对照表集中在 python/pathway/xpacks/llm/constants.py。
4. 检索器工厂:BruteForceKnnFactory
BruteForceKnnFactory 负责构建 KNN(K 近邻)索引,定义在 python/pathway/stdlib/indexing/nearest_neighbors.py。它直接持有上述 embedder,查询时把用户问题也向量化,再在全部文档向量中做暴力最近邻搜索,返回与问题最相似的前 k 个文档块。对中小规模文档集,"暴力全量扫描 + 实时增量维护"在 Pathway 中依然是简单且准确的选择。
5. 组装:DocumentStore
DocumentStore(见 python/pathway/xpacks/llm/document_store.py)是摄取与检索的统一入口。文档进入后依次经过 parser → splitter → embedder → 索引,全链路以 Pathway 数据流形式实时运转。正因为索引构建是增量、实时的,模板 README 特别强调:RAG 默认是"活"的(live)——文档一旦变化,索引随即更新,无需人工重建或重启。
查询处理与向量检索
HTTP 服务器收到的原始查询先经过 prepare_queries(),把它规整成 DocumentStore 检索接口期望的字段:
def prepare_queries(raw_queries: QuerySchema) -> RetrieveQuerySchema:
queries = raw_queries.select(
query=pw.this.messages,
k=1,
metadata_filter=None,
filepath_globpattern=None,
)
return queries
这里把请求体的 messages 映射为检索用的 query,并把 k 固定为 1(即每个问题只召回最相关的 1 个文档块)。RetrieveQuerySchema 与检索结果的结构在 main.py 中单独声明,并且 DocumentStore 自身也内置了同名的 RetrieveQuerySchema / QueryResultSchema 可复用(见 document_store.py)。RetrieveQuerySchema 的字段含义:
query:用于相似度检索的文本;k:需要返回的文档块数量;metadata_filter:按元数据(如文件路径、来源)过滤检索范围,可传None表示不过滤;filepath_globpattern:按文件路径的 glob 通配模式限定候选文档。
实际检索调用发生在 pipeline():
retrieved_documents: QueryResultSchema = document_store.retrieve_query(queries)
retrieved_documents: RetrievedDocsSchema = retrieved_documents.select(
docs=pw.this.result
)
DocumentStore.retrieve_query()(document_store.py 中定义)对每行查询返回与之最接近的文本列表,结果以 JSON 形式存放于 result 字段,随后被包装成 docs 供下游拼装上下文。
上下文拼接与提示词构造
检索出的相关文档要在 build_prompts() 中拼进提示词:
queries_context = queries + documents
def get_context(documents):
content_list = []
for doc in documents:
content_list.append(str(doc["text"]))
return " ".join(content_list)
@pw.udf
def build_prompts_udf(documents, query) -> str:
context = get_context(documents)
prompt = (
f"Given the following documents : \n {context} \nanswer this query: {query}"
)
return prompt
prompts = queries_context + queries_context.select(
prompts=build_prompts_udf(pw.this.docs, pw.this.query)
)
要点解读:
queries + documents是一次表连接(join):查询表与检索结果表按共有列对齐,把每个问题与其召回文档放到同一行;build_prompts_udf被@pw.udf装饰为普通 Python 函数(UDF),Pathway 会在流式数据到达时逐行调用;- 最终拼出的 prompt 形如
Given the following documents : \n <context> \nanswer this query: <query>——即"给定如下文档 + 请回答该问题"。
如果你希望调整问答模板的语气、输出格式(如"只根据文档回答,文档中没有就回答不知道"),替换这一段 prompt 即可;仓库的 40.rag-customization/20.custom-prompt.md 对这一定制做了专门讲解。
基于 LLM Chat 的答案生成
拼好的 prompt 在 get_responses() 中交给大语言模型:
def get_responses(prompts: PromptsSchema) -> ResultsSchema:
model = llms.OpenAIChat(
model="gpt-4o-mini",
api_key=os.environ["OPENAI_API_KEY"],
)
response = prompts.select(
*pw.this.without(pw.this.query, pw.this.prompts, pw.this.docs),
result=model(
llms.prompt_chat_single_qa(pw.this.prompts),
),
)
return response
llms.OpenAIChat是 Pathway 对 OpenAI Chat 服务的封装(见 llms.py),模板选用轻量且成本较低的gpt-4o-mini;llms.prompt_chat_single_qa(...)把一段问题文本包装成 OpenAI chat 所需的消息结构(单条role/content的用户消息),其定义同样位于 llms.py;pw.this.without(...)在输出前剔除query、prompts、docs等中间列,只保留查询标识与最终生成的result;- 对 UDF 形式的 LLM 调用,Pathway 还提供
capacity(并发度)、retry_strategy(重试)、cache_strategy(结果缓存)等参数用于生产调优。
最后,在 pipeline() 中把回答表交给 HTTP writer 回写给调用方,并用 pw.run(monitoring_level=pw.MonitoringLevel.NONE) 启动实时引擎:
writer(responses)
pw.run(monitoring_level=pw.MonitoringLevel.NONE)
monitoring_level=pw.MonitoringLevel.NONE 表示关闭监控埋点,适合最小化运行的示例场景。
运行与查询
第一步:运行流水线
进入 examples/projects/question-answering-rag 目录(并确保已放置好 ./data/ 文档与 .env 配置),执行:
python main.py
启动后,Pathway 引擎会开始摄取文档并建立索引,同时把 HTTP 服务跑在 0.0.0.0:8011。由于索引是增量实时维护的,运行期间若 ./data/ 中的文档被新增、修改或删除,流水线会自动反映这些变化。
第二步:用 curl 发送查询
另开一个终端,向服务器 POST 一条 JSON 消息:
curl --data '{ "messages": "What is the value of X?"}' http://localhost:8011
请求体遵循 QuerySchema 的 messages 字段格式。
第三步:接收回答
服务器处理完查询后,会返回一个基于检索文档生成的回答——回答内容来自"真实文档",而非模型凭空发挥。
注意,rest_connector 的端口与请求格式完全由上面提到的 get_webserver() 决定;如果改了端口或 schema,请求方式也要相应调整。仓库的 40.rag-customization/10.REST-API.md 对 REST API 的扩展用法有更完整的说明。
小结与延伸阅读
本文基于 question-answering-rag 模板说明 及其 main.py,串讲了 Pathway RAG 问答示例的完整结构:用 UnstructuredParser 解析并分块文档 → 用 TokenCountSplitter 精切 token → 用 OpenAIEmbedder 向量化 → 用 BruteForceKnnFactory 建立 KNN 检索 → 由 DocumentStore 统一管理文档摄取与 retrieve_query 检索 → 拼接带上下文的 prompt → 交给 OpenAIChat 生成回答 → 通过 PathwayWebserver + rest_connector 对外提供 HTTP 问答服务。
该示例最值得注意的特性是"RAG 默认是实时的":文档一旦变化,索引随即增量更新,回答始终基于最新的文档内容,无需在每次文档变更后重启或重建索引。这也正是 Pathway 把 RAG 流水线做成流式数据流(而非一次性批处理脚本)的价值所在。
若想继续深入,仓库还提供了可直接扩展的配套材料:
- 更完整的 RAG 问答示例讲解:1000.demo-question-answering.md;
- 更多即开即用的 AI 流水线模板:自适应 RAG 1001.template-adaptive-rag.md、私有化 RAG 1002.template-private-rag.md、多模态 RAG 1003.template-multimodal-rag.md;
- 按需定制各环节(REST API、自定义 prompt、解析器、切分器、embedder、LLM Chat)的专题:docs/2.developers/7.templates/40.rag-customization 目录。
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