pathway
Python ETL 框架,用 Rust 引擎实现批流一体处理,为 LLM/RAG 管道提供实时数据
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
Python ETL 框架,用 Rust 引擎实现批流一体处理,为 LLM/RAG 管道提供实时数据
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
假设你负责维护一个金融信息平台,数据来源包括十几种外部 API、新闻 RSS、交易所 WebSocket 推送以及用户上传的 CSV 文件。传统的处理方式是每天凌晨跑一次批处理任务,把所有数据汇总后生成报表。但问题是:你的 AI 问答机器人回答用户提问时,调用的索引是"昨天凌晨生成"的——用户问今天某只股票的走势,机器人只能尴尬地说"我没有今天的数据"。
Pathway 正是为解决这类问题而生。它是一个 Python 编写的实时数据处理框架,能够让数据管道持续运行在新数据上,每次数据更新都会触发增量计算,让下游的 AI 应用始终"看到"最新状态。
Pathway 由 Pathway 公司开发并维护,定位为 Python ETL 框架,但与传统 ETL 工具(如 Airbyte、Fivetran)不同,它的独特之处在于将批处理与流处理统一在同一个引擎下,且专为 AI/RAG 场景深度优化。截至 2025 年,该项目在 GitHub 上拥有超过 6.3 万颗星,1,600 多次 fork,是数据工程领域增长最快的开源项目之一。
Pathway 的设计哲学是:用纯 Python API 编写数据处理逻辑,但实际计算由底层 Rust 引擎驱动。这意味着开发者可以像写普通 Python 脚本一样写复杂的数据管道,同时获得 Rust 级别的性能和内存效率——多线程、多进程、分布式计算全部开箱即用。
Pathway 的许可证为 BSL 1.1(Business Source License),允许免费用于非商业场景和大多数商业用途,代码在 4 年后自动转为 Apache 2.0 开源许可证。

图1:Pathway 数据处理示例 —— 实时聚合购物列表数据
Pathway 最大的技术亮点是批处理与流处理同引擎。传统架构中,开发者通常需要维护两套代码——批处理脚本(每日/每小时运行)和流处理系统(Kafka + Flink/Spark Streaming)——来处理历史数据和实时数据。而 Pathway 让同一段代码同时支持两种模式:
这种"写一次,跑两种场景"的能力大幅降低了数据管道的维护成本。
Pathway 的底层基于 Differential Dataflow(差分数据流)理论实现,这是一种专为增量计算设计的分布式计算模型。当新的数据点到达时,引擎不会重新计算整个管道,而是只计算受到影响的局部结果。官方提供的 benchmark 显示,在 WordCount 等典型任务上,Pathway 的性能显著优于 Apache Flink、Apache Spark Streaming 和 Kafka Streams。
Pathway 原生支持多种数据源和数据目的地(Connector):
| 类型 | 支持的连接器 |
|---|---|
| 消息队列 | Apache Kafka, Redpanda, Apache Pulsar |
| 数据库 | PostgreSQL, Delta Lake, SQL (via SQLGlot) |
| 云存储 | S3/GCS/Azure Blob, Google Drive, SharePoint |
| 消息格式 | Debezium CDC (MongoDB, PostgreSQL) |
| 文件格式 | CSV, JSON, Parquet, Apache Arrow |
| Airbyte | 300+ 数据源(通过 Airbyte 连接器) |
如果现有连接器无法满足需求,Pathway 提供了 Python Connector 框架,允许用户用纯 Python 编写自定义数据源和数据目的地。
Pathway 单独发布了 xpack-llm 扩展,将大语言模型能力直接集成到数据管道中:
这一套组合使得构建"实时文档 RAG 系统"变得极其简单:Pathway 监控文件系统的更新 → Docling 解析文档内容 → LLM 生成嵌入向量 → 更新向量索引 → AI 应用查询最新数据。

图2:Pathway 官方 GitHub 头像
Pathway 采用混合语言架构:Python 提供用户 API 层,Rust 实现核心计算引擎。项目使用 Maturin 作为构建工具,将 Rust 代码编译为 Python 扩展模块,实现零开销的跨语言调用。
核心目录结构:
python/pathway/ # Python API 封装层
io/ # 数据连接器(CSV、JSON、Kafka、SQL等)
stdlib/ # 标准库(join、groupby、window等)
xpacks/ # 扩展包(llm、sharepoint、milvus等)
udfs.py # 用户自定义函数支持
schema.py # 数据Schema定义
persistence/ # 状态持久化
src/ # Rust 核心引擎( Differential Dataflow 实现)
examples/projects/ # 部署示例(Kafka ETL、RAG、Fargate/Azure部署等)
examples/notebooks/ # Jupyter Notebook 教程
这种架构的优势在于:Python 开发者无需学习 Rust 即可使用 Pathway 的全部能力,同时性能关键的路径完全在 Rust 侧执行。
Pathway 对用户最友好的地方是安装方式——一行 pip 命令即可安装核心包:
pip install pathway
要求 Python 3.10 及以上,支持 macOS 和 Linux。xpack-llm 等扩展模块需要单独安装:
pip install "pathway[xpack-llm]"
pip install "pathway[xpack-llm-local]" # 本地模型(Ollama)
pip install "pathway[xpack-llm-docs]" # 文档解析
Pathway 提供了官方 Docker 镜像 pathwaycom/pathway:latest,开发者只需把自己的处理逻辑和依赖写入 Dockerfile 即可:
FROM pathwaycom/pathway:latest
WORKDIR /app
COPY requirements.txt ./
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
CMD ["python", "./your-script.py"]
也可以直接在宿主机编写 Python 脚本,用 Docker 运行:
docker run -v "$PWD":/app pathwaycom/pathway:latest python my-script.py
官方还提供了 AWS Fargate 和 Azure ACI 的完整部署模板(含 Dockerfile 和启动脚本),可以在云上以无服务器方式运行 Pathway 管道,按实际计算时间计费。
Pathway 的容器化特性使其天然适合 Kubernetes 部署。官方表示企业版支持分布式 Kubernetes 部署和外部持久化配置,适合大规模数据处理场景。
Pathway 内置了一个 Web 监控面板,实时展示各连接器的消息吞吐量、系统延迟和日志信息。开发者无需额外配置监控工具,就能直观了解管道的运行状态。
Pathway 的 API 设计高度贴近 Pandas 和 SQL 的直觉,上手门槛对 Python 开发者非常友好。以下是一个实时处理 CSV 数据流的核心代码示例:
import pathway as pw
class InputSchema(pw.Schema):
value: int
# 读取输入(支持静态CSV或实时Kafka)
input_table = pw.io.csv.read("./input/", schema=InputSchema)
# 数据转换
filtered_table = input_table.filter(input_table.value >= 0)
result_table = filtered_table.reduce(
sum_value = pw.reducers.sum(filtered_table.value)
)
# 输出到文件
pw.io.jsonlines.write(result_table, "output.jsonl")
# 启动管道(批处理直接运行,流处理持续监听)
pw.run()
这段代码的核心结构与 Pandas 无异,但当输入数据更新时,结果会自动增量更新——无需修改任何代码逻辑。
Pathway 提供了 Jupyter Notebook 教程,配合 Google Colab 可以直接在浏览器中体验完整功能。官方还提供了 Cookiecutter 项目模板,帮助开发者快速初始化新项目。
尽管 Pathway 能力强大,但使用时也有需要注意的局限性:
License 限制:BSL 1.1 许可证在商业使用上有一定限制——如果将 Pathway 集成到自己的商业产品中分发,可能需要获得商业许可或等待 4 年许可证自动转为 Apache 2.0。
Python GIL 问题的两面性:虽然 Rust 引擎绕过了 GIL,但用户自定义函数(UDF)仍在 Python 侧执行,对于 CPU 密集型 UDF 需要额外注意性能。
Windows 兼容性:官方仅支持 macOS 和 Linux,Windows 用户需要通过虚拟机或 WSL 运行。
学习曲线:对于完全不了解流处理的数据工程师,Differential Dataflow 的增量计算模型与传统批处理思维方式有较大差异,需要一定时间适应。
在数据工程领域,Pathway 的出现填补了"Python 友好 + 流批一体 + AI 原生"三个需求的交集空白。Apache Kafka + Apache Flink 组合功能强大,但 Java 生态对 Python 开发者不够友好;Apache Spark 有 PySpark 支持,但流处理能力相对较弱;RisingWave 虽然也走 Rust + 流批一体路线,但定位更偏向云原生数据库而非数据管道工具。
Pathway 在 RAG/LLM 应用的数据管道这一细分场景上具有显著优势——它将文档解析、流式更新、向量索引维护和 LLM 调用整合在同一个框架中,解决了 AI 应用最头疼的"数据新鲜度"问题。LangChain 和 LlamaIndex 官方均已将 Pathway 列为推荐的数据连接器。
从 GitHub Stars 增长曲线来看,Pathway 的增长势头在 2024-2025 年明显加速,与 RAG 应用的大规模爆发高度相关。可以预见,随着企业级 RAG 和 AI Agent 应用的普及,Pathway 作为"AI 数据底座"的定位将愈发重要。