库 / SDK
pathwaycom/pathway avatar
pathwaycom/pathway

Pathway:用 Python 写流处理,用 Rust 引擎跑增量计算

Pathway 是面向流处理、实时分析、LLM 流水线与 RAG 的 Python ETL 框架,底层由基于 Differential Dataflow 的可扩展 Rust 引擎驱动。

62,284 个 Star1,685 个 ForkPython许可证因项目而异

秒懂

它是什么?
Pathway 是一个 Python ETL 框架,面向流处理、实时分析和 RAG 管道。它的核心是用 Rust 实现的增量计算引擎,同一份代码可以跑批处理也能跑流式任务。本文基于官方 README 和仓库信息,分析它的机制、适用场景和需要注意的边界。
适合谁用?
Pathway 适合那些已经熟悉 Python、但需要处理实时数据流或构建 RAG 管道的团队。如果你希望用一套代码同时覆盖批处理和流处理,并且愿意接受全内存计算带来的资源约束,Pathway 值得认真评估。
能商用吗?
请先确认。这个仓库使用的许可证不在我们自动归类的范围内,商用前请阅读仓库里的 LICENSE 文件。
还在维护吗?
在维护。仓库在最近一天内有新的提交。
用什么语言写的?
主要是 Python(依据 GitHub 的语言统计)。

以上回答依据项目的 GitHub 数据(最近同步于 2026年9月15日)和我们的分析,不构成法律意见。

开源项目深度解析

它解决什么问题,谁该关注

Pathway 定位是 Python ETL 框架,但它的目标不是传统的批处理调度,而是让开发者用 Python 写出能够持续处理实时数据的管道。官方描述里反复强调两个关键词:流处理和增量计算。传统做法是先用批处理框架处理历史数据,再单独搭一套流处理系统处理新数据,两套代码维护成本高。Pathway 想用一套代码同时覆盖两种模式,让同一个管道既能回放历史数据,也能实时响应新事件。适合的人群包括:需要实时仪表盘的团队、做在线特征计算的机器学习工程师、以及正在搭建 RAG 应用但不想维护多个独立组件的开发者。

Rust 引擎与 Differential Dataflow 的增量机制

Pathway 的底层不是 Python 解释器,而是一个用 Rust 写的执行引擎,基于 Differential Dataflow 做增量计算。这意味着你写的 Python 代码会被翻译成数据流图,然后由 Rust 引擎执行。文档里说这个引擎支持多线程、多进程和分布式计算,而且整个管道的数据都保存在内存中。增量计算的核心好处是:当新数据到达时,引擎只重新计算受影响的部分,而不是从头跑一遍整个管道。这对流式场景很关键,因为数据是持续进来的,如果每次都要全量重算,延迟和资源消耗都不可接受。但这也带来一个隐含约束:全内存设计意味着数据规模受限于集群内存总量,这一点在规划部署时就要想清楚。

安装与第一个管道:从 pip 到运行

安装很简单,官方文档给出的命令是 pip install -U pathway,要求 Python 3.10 或更高版本。安装之后,你可以用 Python API 定义数据源、转换逻辑和输出。README 里没有给出完整的代码示例,但提到了连接器的概念,比如 Kafka、GDrive、PostgreSQL 和 SharePoint。一个典型的管道可能从 Kafka 读取事件流,用内置的转换函数做窗口聚合,然后写入数据库或触发告警。官方提供了多个可运行的模板,包括 realtime-log-monitoring 和 linear_regression_with_kafka,这些模板同时提供 notebook 和 Docker 格式,可以直接启动。对于想快速验证的开发者,先跑一个模板比从零写代码更实际。

批流统一:同一份代码的两种运行方式

Pathway 的一个核心卖点是批流统一。文档说同一份代码可以用于本地开发、CI/CD 测试、批处理作业、流重放和实时流处理。这意味着你可以在开发环境用批模式调试逻辑,然后部署到生产环境切换成流模式,不需要改代码。这种设计降低了从批处理迁移到流处理的成本,也方便做回测:用历史数据重放管道,验证逻辑正确性。但要注意,批流统一不等于没有差异。流处理中的时间语义、迟到数据和乱序事件在批处理中不存在,Pathway 声称会处理这些问题,但你的代码仍然需要理解这些概念,否则可能写出在批模式下正确、在流模式下出错的处理逻辑。

一致性模型:免费版与商业版的差异

文档明确区分了两种一致性级别:免费版提供 at least once 一致性,企业版提供 exactly once 一致性。这是一个重要的边界。at least once 意味着在故障或重启时,某些数据可能被重复处理,这可能导致重复写入或重复计算。对于某些场景,比如日志监控或趋势分析,重复一条数据影响不大。但对于金融交易或库存扣减,重复处理可能造成严重问题。如果你需要 exactly once,就必须考虑企业版,而这会带来额外的成本和商业条款约束。在选型时,先问自己:我的业务能容忍重复吗?如果不能,免费版可能不够用。

持久化与重启恢复:状态保存的代价

Pathway 提供持久化功能,可以保存计算状态,以便在更新或崩溃后重启管道。这是流处理系统的关键能力,因为没有状态保存,重启就意味着丢失所有中间结果,需要从头消费数据。文档说持久化允许你重启管道,但没有详细说明状态存储的格式或恢复的具体机制。从设计上推断,状态保存在内存中,持久化可能是把状态快照到外部存储。这里有一个权衡:持久化会增加写入开销,频繁的快照可能影响性能,而快照间隔太长又可能丢失较多状态。官方文档没有给出具体的性能数字,所以你需要根据自己的数据量和恢复时间要求做测试。

LLM 与 RAG 的集成:不只是 ETL

Pathway 不只是做传统 ETL,它还提供了专门的 LLM 扩展包,包含 LLM 包装器、解析器、嵌入器和分割器,以及一个内存中的实时向量索引。这意味着你可以直接在 Pathway 里构建 RAG 管道,而不必把数据导出到另一个向量数据库。文档提到了与 LlamaIndex 和 LangChain 的集成,这降低了已有 LLM 应用的迁移成本。官方模板包括 private-rag-ollama-mistral 和 adaptive-rag,这些模板展示了如何用 Pathway 处理实时文档更新并保持索引同步。对于做 RAG 的团队来说,Pathway 的价值在于把数据摄入、转换和向量检索放在同一个管道里,减少了系统组件数量。但要注意,向量索引在内存中,所以文档规模受限于内存容量,大规模知识库可能需要分片或考虑其他方案。

替代方案与选型考量

Pathway 的主要替代方案是 Apache Flink 和 Kafka Streams。Flink 也提供流处理和有状态计算,但它的 API 是 Java/Scala 为主,Python API 支持相对有限。Kafka Streams 是库而不是框架,它嵌入在你的应用中,适合已经重度使用 Kafka 的团队。与这些相比,Pathway 的优势是纯 Python API,降低了学习门槛,而且内置了 LLM 相关工具,这是 Flink 和 Kafka Streams 没有的。但 Flink 和 Kafka Streams 在分布式部署和状态管理方面更成熟,社区更大,文档更丰富。Pathway 的全内存设计可能限制它的扩展性,而 Flink 可以 spill 到磁盘。如果你已经有 Kafka 基础设施,Kafka Streams 可能更自然;如果你需要 Python 生态和 LLM 集成,Pathway 更合适。

编辑结论

Pathway 适合那些已经熟悉 Python、但需要处理实时数据流或构建 RAG 管道的团队。如果你希望用一套代码同时覆盖批处理和流处理,并且愿意接受全内存计算带来的资源约束,Pathway 值得认真评估。但如果你只需要简单的批处理 ETL,或者对数据一致性要求极高且无法接受至少一次语义,那么它可能不是首选。在采用之前,先确认三件事:你的数据源是否在官方连接器列表内,你的管道能否容忍重启后的状态恢复延迟,以及你需要的特性是否落在免费版和商业版的功能分界线上。许可证文件在仓库根目录的 LICENSE.txt 中,部署前请仔细阅读,特别是商业使用条款。

官方来源

  1. Official documentation
  2. Official README
  3. Project repository
  4. Release notes
社区笔记

社区笔记