首页
/ Pathway 实战:用 Pathway 构建可实时更新的检索增强生成(RAG)问答流水线

Pathway 实战:用 Pathway 构建可实时更新的检索增强生成(RAG)问答流水线

2026-09-07 09:39:41作者:魏侃纯Zoe

导读

本文围绕仓库中的 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 目录,通过"检索 + 生成"的组合提升回答的准确性与相关性:先召回与用户问题最相关的文档片段,再让大语言模型基于这些片段生成答案。这样既保证了生成的连贯性,也让回答内容有据可依、贴近最新文档。

模板由两部分组成:

  1. RAG 流水线本体main.py):负责文档摄取、向量索引、查询检索与答案生成;
  2. 轻量 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 组件(DocumentStoreUnstructuredParser、各类 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.pyOpenAIChat 的构造说明)。

数据准备与 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.pyUnstructuredParser 的定义。

chunking_kwargs 中的 max_characters=3000new_after_n_chars=2000 用于限制单块字符数与换块时机,防止切出过长的文本块。

2. 切分器:TokenCountSplitter

即便解析器已经按标题切块,模板仍用 TokenCountSplitter 按 token 数做二次精切:min_tokens=100max_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(...) 在输出前剔除 querypromptsdocs 等中间列,只保留查询标识与最终生成的 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

请求体遵循 QuerySchemamessages 字段格式。

第三步:接收回答

服务器处理完查询后,会返回一个基于检索文档生成的回答——回答内容来自"真实文档",而非模型凭空发挥。

注意,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 流水线做成流式数据流(而非一次性批处理脚本)的价值所在。

若想继续深入,仓库还提供了可直接扩展的配套材料:

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

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388