Pathway LLM/RAG 实战示例详解:用 LLM xpack 构建实时 RAG、Adaptive RAG 与文档索引流水线
本文围绕 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 Indexing 与 Vector Store 说明。
源码印证:DocumentStore 到底做了什么
在 DocumentStore 实现中可以看到它的完整职责:读入含 data(bytes 列,通常由连接器 format="raw"/"binary" 产生)与可选 _metadata 列的 Table,用 parser 解析、splitter 切块、doc_post_processors 可选后处理,再由 retriever_factory 构建向量索引并暴露 retrieve_query 等查询方法。
查询中的过滤能力同样可在源码中验证:metadata_filter 与 filepath_globpattern 会先经 _get_jmespath_filter 合并成一条 JMESPath 过滤表达式(路径过滤使用 globmatch 对 path 字段做通配匹配),这解释了为什么官方示例中这两个字段可以直接传 None。DocumentStore 还实现了 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 实现,它继承自 BaseRAGQuestionAnswerer(L442),后者又继承自 SummaryQuestionAnswerer → BaseQuestionAnswerer(L388)。也就是说,模板 YAML 中只需声明 BaseRAGQuestionAnswerer 即可得到标准 RAG;换成 AdaptiveRAGQuestionAnswerer 就获得动态检索行为。相关行为有专门的测试覆盖,见 test_rag.py。
精选示例三:全私有 RAG(Mistral + Ollama)
private_rag 示例展示如何搭建完全不依赖外部云服务的私有 RAG:LLM 与嵌入模型都跑在本地(通过 Ollama 提供 Mistral 等模型),检索仍使用 Pathway 的自适应检索。
- 模板文档:Private RAG App with Mistral and Ollama(对应
templates/private_rag) - 深入文章:Private RAG with Connected Data Sources using Mistral, Ollama, and Pathway
从源码结构看,本地化替换的本质是换 LLM 封装与 Embedder:在模板 YAML 中把 llm 换成走本地 api_base 的 LiteLLMChat、把 embedder 换成 SentenceTransformerEmbedder,整条流水线就不再调用任何外部服务——这正是 YAML 配置指南中给出的私有 RAG 构造方法。可用封装列表见 LLM Chats 文档 与 Embedders 文档。
精选示例四:多模态 RAG(GPT-4o 解析财务文档)
multimodal_rag 示例针对一个传统 RAG 的明显短板:基于表格/图表数据的问答。它把文档处理流水线中解析环节交给 GPT-4o 这类多模态模型,从文件夹里的非结构化财务文档(含表格、图表等视觉元素)中提取信息;文档变化或新增时流水线自动更新结果,使 AI 应用与文档驱动器保持常驻连接。
- 模板文档:Multimodal RAG with Pathway(对应
templates/multimodal_rag) - 案例讲解:Multimodal RAG 文章
- 仓库内还有可直接运行的演示 notebook:multimodal-rag.ipynb 与 multimodal-rag-using-Gemini.ipynb
其他示例:文档索引、Drive 告警与非结构化转 SQL
Realtime Document Indexing(document_indexing)
基础版实时文档索引流水线:从 SharePoint / Google Drive 等数据源索引文档,之后可以查询索引、获取索引统计、读取文件元数据。
- 模板文档:Document Indexing 模板
- 原理与手写实现:Document Indexing 教程
- 演示 notebook:live_vector_indexing_pipeline.ipynb
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 问答流水线(含自适应检索变体)与 RAGClient(L1070) |
LLM xpack 总览 |
| Reranker | rerankers.py | 对检索结果重排以提升相关性 | Rerankers 文档 |
组件行为的正确性由 xpack 测试套件持续验证,例如 test_document_store.py、test_rag.py、test_parsers.py、test_splitters.py 与 test_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_rag 与 private_rag 分别演示了动态检索降本与全本地化部署;multimodal_rag 补齐了表格/图表类文档的短板;而 document_indexing、drive_alert、unstructured_to_sql_on_the_fly 则覆盖了索引、事件告警与非结构化数据落库三类高频场景。配合 LLM xpack 总览、xpack 源码 与 集成测试,你既可以照抄模板快速上线,也可以下钻到 DocumentStore、AdaptiveRAGQuestionAnswerer 等实现细节,按需定制自己的实时 LLM 流水线。
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 StartedRust0623
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