首页
/ Pathway LLM/RAG 实战示例详解:用 LLM xpack 构建实时 RAG、Adaptive RAG 与文档索引流水线

Pathway LLM/RAG 实战示例详解:用 LLM xpack 构建实时 RAG、Adaptive RAG 与文档索引流水线

2026-09-06 09:07:20作者:韦蓉瑛

本文围绕 Pathway(Pathway Live Data Framework)官方整理的 LLM 示例集合(LLM Examples)展开:它用一组可直接运行的 RAG / 文档索引 / 非结构化数据提取示例,展示了如何用 Pathway 的 LLM 工具链(pathway.xpacks.llm)在无需独立 ETL 的前提下构建“知识永远最新”的大模型应用。读完后,你将掌握每个示例对应的实现原理——从 DocumentStore 的索引机制、UnstructuredParser/TokenCountSplitter 的文档处理,到 AdaptiveRAGQuestionAnswerer 的动态检索策略,并能在本仓库中找到每个示例对应的教程文档与源码实现位置。

示例总览:Pathway LLM 工具链能做什么

Pathway 的 LLM xpack(pathway[xpack-llm])为流式数据处理框架内置了一整套 LLM/RAG 组件:文档解析器(parsers)、文本切分器(splitters)、向量化器(embedders)、重排器(rerankers)、LLM 聊天封装(llms)、预置 Prompt(prompts)、HTTP 服务封装(servers)以及整条问答流水线(question answering)。官方 LLM 示例页正是基于这些组件收集的一组实战模板,全部以 llm-app 模板仓库中的 templates/ 目录形式提供,并支持 Docker 一键运行。

官方示例分为两类:Featured examples(精选示例)与 Other examples(其他示例)。下表完整继承自原示例页:

示例 说明 本仓库对应文档
question_answering_rag:带“永远最新知识”的 RAG 应用 展示如何用 Pathway 创建 RAG 应用,无需单独 ETL 即可为 LLM 提供始终最新的知识。可从 SharePoint、Google Drive 等不同数据源构建向量库,并用任意 LLM 模型基于索引文档回答问题 Create your own RAG
adaptive_rag:Adaptive RAG 展示 Adaptive RAG 技术:利用 LLM 的反馈动态调整 RAG Prompt 中携带的文档数量 Adaptive RAG App
private_rag:全私有 RAG 展示如何用 Pathway、Mistral 与 Ollama 搭建带自适应检索的私有 RAG 流水线 Private RAG App with Mistral and Ollama
multimodal_rag:多模态 RAG 启动一个依赖文档处理流水线的多模态 RAG,解析阶段使用 GPT-4o。Pathway 从文件夹中的非结构化财务文档提取信息,文档变更或新文档到达时自动更新结果,使 AI 应用与文档驱动器保持常驻连接、实时同步,并支持表格、图表等视觉元素;相比传统 RAG 难以回答表格类问题的局限,多模态方案在表格信息抽取上表现更佳 Multimodal RAG 模板
document_indexing:实时文档索引 一个基础的实时文档索引流水线示例:从 SharePoint、Google Drive 等不同数据源索引文档,可查询索引、获取索引统计信息与文件元数据 Document Indexing 模板
drive_alert:Drive 告警流水线 与 “Alert” 示例几乎相同,唯一区别是数据源换成 Google Drive Drive Alert 模板
unstructured_to_sql_on_the_fly:非结构化转 SQL 从非结构化数据(PDF 与查询语句)中“即时”抽取并结构化数据 Unstructured to SQL

精选示例一:知识永远最新的 RAG 问答应用

question_answering_rag 是 LLM 示例中的核心模板。它的核心价值主张是:索引实时化是 RAG 答案不“过期”的唯一途径。Pathway 的 Live Data Framework 把静态数据与流式数据用同一套模型处理,因此文档目录一旦发生变化(新增、修改、删除文件),索引会自动增量更新,RAG 的回答始终基于最新文档。

安装与前置条件

Create your own RAG 教程,先安装带 LLM xpack 的 Pathway:

pip install pathway[xpack-llm] python-dotenv

并将 OpenAI API key 放入 .env 文件:

OPENAI_API_KEY="sk-..."

完整流水线代码骨架

RAG 流水线的七个环节为:文档索引 → 用户查询 → 文档检索 → 上下文构建 → Prompt 构造 → 答案生成 → 返回输出。下面按官方教程给出完整可运行骨架(以本地 ./data/ 目录中的 PDF 为例)。

1. 导入 LLM xpack 组件:

import pathway as pw
from pathway.stdlib.indexing.nearest_neighbors import BruteForceKnnFactory
from pathway.xpacks.llm import llms
from pathway.xpacks.llm.document_store import DocumentStore
from pathway.xpacks.llm.embedders import OpenAIEmbedder
from pathway.xpacks.llm.parsers import UnstructuredParser
from pathway.xpacks.llm.splitters import TokenCountSplitter

from dotenv import load_dotenv
import os

load_dotenv()

2. 通过文件系统连接器读入文档:

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

3. 组装 DocumentStore 的四个核心部件——文本切分器(按 token 数切块)、向量化器(把文本变成语义向量)、检索器工厂(基于向量找最相关文档)、解析器(从文档中提取并结构化文本):

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,
    },
)

document_store = DocumentStore(
    docs=documents,
    retriever_factory=retriever_factory,
    parser=parser,
    splitter=text_splitter,
)

4. 用轻量 HTTP 服务接收用户查询

webserver = pw.io.http.PathwayWebserver(host="0.0.0.0", port=8011)

class QuerySchema(pw.Schema):
    messages: str

queries, writer = pw.io.http.rest_connector(
    webserver=webserver,
    schema=QuerySchema,
    autocommit_duration_ms=50,
    delete_completed_queries=False,
)

5. 把查询整理成 DocumentStore.retrieve_query 期望的格式,字段包括 query(用户问题)、k(检索文档数)、metadata_filter(可选,按元数据过滤)、filepath_globpattern(可选,按路径模式过滤文件)。这里为所有查询统一取 k=1、不做过滤;也可在 QuerySchema 中增加 k 字段实现按查询定制:

queries = queries.select(
    query = pw.this.messages,
    k = 1,
    metadata_filter = None,
    filepath_globpattern = None,
)

6. 检索、构建上下文与 Prompt:

retrieved_documents = document_store.retrieve_query(queries)
retrieved_documents = retrieved_documents.select(docs=pw.this.result)
queries_context = queries + retrieved_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)
)

7. 用 LLM 生成答案并写回:

model = llms.OpenAIChat(
    model="gpt-4o-mini",
    api_key=os.environ["OPENAI_API_KEY"],  # 从环境变量读取 OpenAI API key
)

responses = prompts.select(
    *pw.this.without(pw.this.query, pw.this.prompts, pw.this.docs),
    result=model(
        llms.prompt_chat_single_qa(pw.this.prompts),
    ),
)

writer(responses)

8. 运行并用 curl 验证:

pw.run()          # 在 main.py 末尾
python main.py
curl --data '{ "messages": "What is the value of X?"}' http://localhost:8011

需要 UI 时可再叠加 Streamlit 等前端。更多索引细节可参考 Document IndexingVector Store 说明

源码印证:DocumentStore 到底做了什么

DocumentStore 实现中可以看到它的完整职责:读入含 data(bytes 列,通常由连接器 format="raw"/"binary" 产生)与可选 _metadata 列的 Table,用 parser 解析、splitter 切块、doc_post_processors 可选后处理,再由 retriever_factory 构建向量索引并暴露 retrieve_query 等查询方法。

查询中的过滤能力同样可在源码中验证:metadata_filterfilepath_globpattern 会先经 _get_jmespath_filter 合并成一条 JMESPath 过滤表达式(路径过滤使用 globmatchpath 字段做通配匹配),这解释了为什么官方示例中这两个字段可以直接传 NoneDocumentStore 还实现了 McpServable 接口(见 MCP Server 文档),因此索引本身可以对外暴露为 MCP 服务。

精选示例二:Adaptive RAG —— 动态控制检索数量

Adaptive RAG 是 Pathway 提出的一种 RAG 优化技术:不再固定检索 k 个文档塞进 Prompt,而是根据 LLM 的反馈动态决定本轮需要多少文档。模型先基于少量上下文作答并给出“证据充分/不充分”的信号,不充分时再扩大检索范围,从而在保证准确率的同时显著降低 token 开销。

  • 模板文档:Adaptive RAG App(对应 llm-app 的 templates/adaptive_rag 目录,支持 Docker 运行)
  • 交互式原理讲解:Adaptive RAG 图文笔记

在源码层面,该技术由 AdaptiveRAGQuestionAnswerer 实现,它继承自 BaseRAGQuestionAnswererL442),后者又继承自 SummaryQuestionAnswererBaseQuestionAnswererL388)。也就是说,模板 YAML 中只需声明 BaseRAGQuestionAnswerer 即可得到标准 RAG;换成 AdaptiveRAGQuestionAnswerer 就获得动态检索行为。相关行为有专门的测试覆盖,见 test_rag.py

精选示例三:全私有 RAG(Mistral + Ollama)

private_rag 示例展示如何搭建完全不依赖外部云服务的私有 RAG:LLM 与嵌入模型都跑在本地(通过 Ollama 提供 Mistral 等模型),检索仍使用 Pathway 的自适应检索。

从源码结构看,本地化替换的本质是换 LLM 封装与 Embedder:在模板 YAML 中把 llm 换成走本地 api_baseLiteLLMChat、把 embedder 换成 SentenceTransformerEmbedder,整条流水线就不再调用任何外部服务——这正是 YAML 配置指南中给出的私有 RAG 构造方法。可用封装列表见 LLM Chats 文档Embedders 文档

精选示例四:多模态 RAG(GPT-4o 解析财务文档)

multimodal_rag 示例针对一个传统 RAG 的明显短板:基于表格/图表数据的问答。它把文档处理流水线中解析环节交给 GPT-4o 这类多模态模型,从文件夹里的非结构化财务文档(含表格、图表等视觉元素)中提取信息;文档变化或新增时流水线自动更新结果,使 AI 应用与文档驱动器保持常驻连接。

其他示例:文档索引、Drive 告警与非结构化转 SQL

Realtime Document Indexing(document_indexing)

基础版实时文档索引流水线:从 SharePoint / Google Drive 等数据源索引文档,之后可以查询索引、获取索引统计、读取文件元数据。

Drive Alert Pipeline(drive_alert)

与 Pathway 官方的 “Alert” 示例几乎一致,唯一区别是数据源改为 Google Drive(监控 Drive 中文件变化并触发告警)。该示例也出现在 Slack 告警连接器文档中,可配合 pw.io.slack.send_alerts 把事件推送到 Slack。

Unstructured to SQL on the Fly

从 PDF 等非结构化数据中即时抽取、结构化数据并写入 SQL。文档展示了整体架构:解析 → 用 LLM 抽取字段 → 落入结构化表。运行方式就是简单的 python app.py,也支持 Docker。详见 Unstructured to SQL 文章

所有示例共享的底层组件:从源码理解 LLM xpack

上述示例并非彼此独立的脚本,而是同一套 xpack 组件的不同组合。结合本仓库源码,可以清楚看到“示例—组件”的映射关系:

组件 源码位置 作用 官方文档
DocumentStore document_store.py 解析→切块→向量化→建索引→retrieve_query 一站式封装 Document Store 说明
UnstructuredParser parsers.py 底层调用 Unstructured 库解析 PDF/Office 等,支持 single/paged/elements 模式与 by_title 等切块策略 Parsers 文档
TokenCountSplitter splitters.py 按 token 数(min_tokens/max_tokens + tiktoken encoding_name)切块,尽量不在句子中间断开 Splitters 文档
OpenAIEmbedder embedders.py 文本→语义向量,OpenAI/SentenceTransformer 等实现 Embedders 文档
llms.OpenAIChat / LiteLLMChat llms.py LLM 提供商封装,实例本身就是可作用于 Table 列的 UDF LLM Chats 文档
BaseRAGQuestionAnswerer / AdaptiveRAGQuestionAnswerer question_answering.py 预置的完整 RAG 问答流水线(含自适应检索变体)与 RAGClientL1070 LLM xpack 总览
Reranker rerankers.py 对检索结果重排以提升相关性 Rerankers 文档

组件行为的正确性由 xpack 测试套件持续验证,例如 test_document_store.pytest_rag.pytest_parsers.pytest_splitters.pytest_embedders.py

如何运行这些示例:模板仓库与 YAML 配置

所有示例统一托管在独立的 llm-app 模板仓库的 templates/ 目录下,每个模板都是 app.py + app.yaml 的组合。本仓库提供了完整的运行指南:

  • Run a template:克隆 llm-app 仓库后进入对应模板目录(如 question_answering_rag),按 README 配置环境变量即可运行,也提供 Docker 方式;
  • Configure YAML:详解 app.yaml$sources$document_store$answerer 等字段的组合方式,以及用大写标识符引用环境变量的语法;
  • YAML 片段示例:以 Adaptive RAG 为例给出可直接复制的完整 YAML 流水线配置。

一个值得注意的事实:question_answering_rag、adaptive_rag、multimodal_rag、private_rag 这几个模板共享同一份 Python 代码,差异只体现在 YAML 配置(见 YAML 配置指南)——这本身就是 Pathway “声明式定义流水线”理念的直观体现:切换 LLM 提供商、切换检索策略、切换文档源,多数情况下只需要改配置。

小结

Pathway 官方 LLM 示例页的价值在于把“LLM xpack 组件”与“可运行模板”连接了起来:question_answering_rag 展示了从零搭建实时 RAG 的完整代码路径;adaptive_ragprivate_rag 分别演示了动态检索降本与全本地化部署;multimodal_rag 补齐了表格/图表类文档的短板;而 document_indexingdrive_alertunstructured_to_sql_on_the_fly 则覆盖了索引、事件告警与非结构化数据落库三类高频场景。配合 LLM xpack 总览xpack 源码集成测试,你既可以照抄模板快速上线,也可以下钻到 DocumentStoreAdaptiveRAGQuestionAnswerer 等实现细节,按需定制自己的实时 LLM 流水线。

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