首页
/ Pathway Live Data Framework:基于 Rust 差分计算引擎的 Python 实时流批一体 ETL 框架

Pathway Live Data Framework:基于 Rust 差分计算引擎的 Python 实时流批一体 ETL 框架

2026-09-03 15:52:39作者:乔或婵

本文以仓库根目录的 README.md 为骨架,系统讲解 Pathway(Pathway Live Data Framework)的定位、核心概念与上手路径,并结合仓库中的 pyproject.tomlCargo.tomlpython/pathway 源码,深入剖析 pw.run() 运行参数、连接器体系、LLM 扩展、持久化机制以及本地/Docker/Kubernetes 三种部署方式,帮助读者掌握这套“Python 写逻辑、Rust 跑引擎”的流式数据框架的完整使用方法。

一、什么是 Pathway:流批一体的 Python ETL 框架

README.md 对 Pathway 的定义是:一个面向流处理(stream processing)、实时分析(real-time analytics)、LLM 管线和 RAG 的 Python ETL 框架。它的设计目标可以归纳为三点:

  1. 易用:提供易用的 Python API,可以无缝集成常用的 Python ML 库;
  2. 同一份代码跨越开发与生产:无论是本地开发、CI/CD 测试、批处理作业、流回放(stream replay)还是实时流处理,都可以使用同一段管线代码;
  3. 性能由 Rust 引擎兜底:引擎基于 Differential Dataflow 实现增量计算(incremental computation),Python 代码由 Rust 引擎执行,从而突破 Python 单线程限制,支持多线程、多进程乃至分布式计算,整个管线驻留在内存中,可轻松用 Docker/Kubernetes 部署。

从源码结构看,这一“Python + Rust 双栈”架构在仓库中有非常清晰的对应:

  • Python 侧位于 python/pathway 目录,提供 API、连接器、LLM 扩展与 CLI;python/pathway/engine.pyi 是 Rust 引擎暴露给 Python 的类型存根;
  • Rust 侧位于 src 目录,其中 src/engine 是计算引擎(dataflow、reduce、frontier 时间推进等),src/connectors 是连接器底层实现,src/persistence 是状态持久化;
  • 计算内核依赖通过本地路径引入:Cargo.tomldifferential-dataflow = { path = "./external/differential-dataflow" }Cargo.tomltimely = { path = "./external/timely-dataflow/timely" },即仓库在 external/differential-dataflowexternal/timely-dataflow 内置了这两个核心 crate 的本地副本;
  • Python 与 Rust 之间通过 PyO3 桥接,Cargo.toml 显示 pyo3 启用了 abi3-py310 特性,与 Python 侧最低 3.10 的版本要求一致;打包由 maturin 完成(见 pyproject.tomlpyproject.toml[tool.maturin] 配置,Rust 库名为 pathway_engine,最终映射为 Python 包内的 pathway.engine 模块)。

Cargo.toml 可以看到当前仓库版本为 0.32.1,要求 Rust 1.97,许可证为 BUSL-1.1

二、安装与环境要求

README.md 的 Installation 一节给出的前提与命令如下:

  • Python 版本要求 3.10 或更高,这与 pyproject.tomlrequires-python = ">=3.10" 完全一致;
  • 平台支持 MacOS 与 Linux,其他系统需通过虚拟机运行(README 原文带 ⚠️ 提示)。

安装命令:

pip install -U pathway

除了核心依赖,pyproject.toml 还声明了 pandasnumpypyarrowfastapiuvicorn、OpenTelemetry 等基础依赖,并提供了按功能拆分的可选依赖组(extras),常用的几组包括:

安装组 用途 关键依赖(摘自源码)
pathway[sql] pw.sql 字符串 SQL 表达式 sqlglot==10.6.1(新版与当前 pw.sql 实现不兼容,故精确锁定)
pathway[xpack-llm] LLM 管线扩展 openailitellmlangchainllama-index-core
pathway[xpack-llm-local] 本地推理 sentence_transformerstransformers
pathway[xpack-llm-docs] 文档解析 doclingunstructuredpaddleocr
pathway[milvus] Milvus 向量库连接器 pymilvusmilvus-lite
pathway[all] 聚合以上全部 pyproject.toml

此外,pyproject.toml 通过 [project.scripts] 注册了 pathway 命令行入口(指向 pathway.cli:main),也就是后文 pathway spawn 命令的来源。

三、快速上手:实时累加正数值的完整示例

下面是 README.md 官方示例(计算正数值之和并实时写出)的完整代码,配合逐段说明:

import pathway as pw

# 1. 定义数据的 Schema(可选,但推荐)
class InputSchema(pw.Schema):
  value: int

# 2. 用连接器接入数据(这里读取 ./input/ 目录下的 CSV)
input_table = pw.io.csv.read(
  "./input/",
  schema=InputSchema
)

# 3. 定义对数据的操作:过滤 + 聚合
filtered_table = input_table.filter(input_table.value >= 0)
result_table = filtered_table.reduce(
  sum_value = pw.reducers.sum(filtered_table.value)
)

# 4. 将结果加载到外部系统
pw.io.jsonlines.write(result_table, "output.jsonl")

# 5. 启动计算
pw.run()

各步骤对应的仓库实现:

除 README 示例外,仓库在 examples/notebooks 下提供 45 个可运行的 Jupyter Notebook(如 examples/notebooks/pathway_intro.ipynb),examples/projects 下还有完整项目样例;官方模板与更多示例可参考 README 指引的在线模板库(notebook 与 docker 两种形态)。

四、pw.run() 运行参数详解

README.md 指出“一行 pw.run() 即可启动流式计算”。结合源码 python/pathway/internals/run.pyrun 的完整签名与参数如下(全部为关键字参数):

参数 默认值 说明
debug False 是否启用 table.debug() 算子的输出
monitoring_level MonitoringLevel.AUTO 监控粒度:NONE / IN_OUT / ALLAUTO 时框架根据输出端是否交互自动在 NONE 与 IN_OUT 间选择
with_http_server False 是否启动带运行时指标的 HTTP 服务(源码注释标注“即将弃用”,监控已演进为内建 dashboard)
default_logging True 是否允许 Pathway 自行设置 logging handler;想接管日志时置为 False
persistence_config None 持久化配置,用于在更新或崩溃后重启管线时恢复计算状态(对应 python/pathway/persistencesrc/persistence
runtime_typechecking None 是否启用更严格的运行时类型检查
terminate_on_error None 数据/用户逻辑出错时是否终止整个计算
max_expression_batch_size 1024 表达式一次计算的最大行数;中间状态较大的表达式可考虑调小
event_loop None 外部创建的 asyncio 事件循环;多次运行含 InMemoryCache 的异步 UDF 时必须复用同一事件循环
udf_cache_directory None 非确定性 UDF 的记忆化缓存目录:设置后缓存以 SQLite 文件存放,避免随缓存条目数增长内存;文件是运行期工作集,每次运行重建、关闭时清理

run 内部构造 GraphRunner(解析全局计算图 parse_graph.G)后调用 run_outputs()python/pathway/internals/run.py);同文件还定义了 run_allpython/pathway/internals/run.py),功能相同但禁用 tree-shaking 优化,即不做“无用输出剪枝”,适合需要完整执行所有注册输出的调试场景。

五、核心特性逐项解析

README.md 的 Features 一节列出了六大特性,这里结合仓库证据逐项展开:

5.1 丰富的连接器生态

README 声明内置 Kafka、GDrive、PostgreSQL、SharePoint 等连接器,且 Airbyte 连接器可接入 300+ 数据源;缺失的连接器可用 Python 自定义。

从源码结构看,python/pathway/io 下的连接器目录规模远超 README 的举例,按字母序包括:airbytebigquerychromaclickhousecsvdebeziumdeltalakeduckdbdynamodbelasticsearchgdriveicebergjsonlineskafkakinesisleannlogstashmilvusminiomongodbmqttmssqlmysqlnatspineconeplaintextpostgrespubsubpulsarqdrantquestdbrabbitmqredpandas3slacksqliteweaviate 等 40 余个子目录。Rust 端 src/connectors/data_storage 则实现了 AWS(S3)、数据湖(Delta Lake/Iceberg)等存储接入;Cargo.toml 的依赖列表也印证了这些底层客户端(rdkafkapg_walstreamicebergelasticsearchquestdb-rsrumqttcpulsar 等)都被直接编进引擎。仓库还在 integration_tests 下按连接器组织了对应集成测试(kafka/db_connectors/s3/airbyte/ 等),可用来验证各连接器的真实读写行为。

5.2 无状态与有状态变换

Pathway 支持 join、窗口(windowing)、排序等有状态变换,且很多变换直接在 Rust 中实现以提升性能;同时任何 Python 函数都可以作为 UDF 参与计算——可以自写,也可以直接调用任意 Python 库。从 src/engine 的模块划分(dataflow/reduce.rsexpression.rs 等)可以看出,聚合、join 等状态算子确实由引擎原生实现,而 Python UDF 则通过 pyo3 回调执行。

5.3 持久化(Persistence)

README 强调持久化可保存计算状态,使管线在更新或崩溃后能够重启恢复。实现上,Python 侧 pw.run(persistence_config=...) 接收 python/pathway/persistence 中的 Config;Rust 侧 src/persistence 包含 state.rsoperator_snapshot.rsfrontier.rs(时间前沿快照)、tracker.rs 以及 backends/(本地文件等后端),与 run.py 中“持久化用于状态保存”的文档说明相互印证。

5.4 一致性(时间管理)

README 说明 Pathway 负责管理时间、保证计算一致性:迟到(late)和乱序(out-of-order)的数据到来时,系统会自动更新既有结果。版本差异上,免费版提供 at-least-once 一致性,企业版提供 exactly-once。这一“以结果增量修正代替重放”的模型正来自差分数据流的 +1/-1 计数语义,与 external/differential-dataflow 的设计一致。

5.5 可扩展的 Rust 引擎

多线程、多进程、分布式计算由引擎负责。Cargo.toml 引入 rayon 做数据并行,tokio 提供异步运行时(src/async_runtime.rs);src/python_api/threads.rs 则对应 Python 侧的多线程控制(即 CLI --threads 参数背后的实现)。

5.6 LLM 工具链

README 声明提供 LLM 扩展:LLM 封装、解析器、embedding、切分器、内存实时向量索引,以及与 LlamaIndex、LangChain 的集成。仓库中的对应证据:

  • python/pathway/xpacks/llm 目录承载 LLM 扩展代码;
  • pyproject.tomlxpack-llm / xpack-llm-local 依赖组(openailitellmcoherelangchainllama-index-coresentence_transformers 等);
  • Rust 端 src/external_integration 实现了内存向量检索(brute_force_knn_integration.rsusearch_integration.rs)与全文检索(tantivy_integration.rsqdrant_integration.rs),Cargo.toml 中的 usearchtantivy 依赖即服务于此;
  • integration_tests/rag_evals 提供了 RAG 管线的评测集成测试。

六、部署方式:本地、Docker、Kubernetes

这一节完整继承 README.md Deployment 章节的操作细节。

6.1 本地运行

创建管线后一行 pw.run() 启动;项目可直接当普通 Python 脚本运行:

$ python main.py

README 说明 Pathway 自带监控 dashboard,可观察每个连接器发出的消息数、系统延迟与日志。注意:本地默认单进程多线程,pathway spawn 可以显式指定线程数:

# 基础方式
$ pathway spawn python main.py

# 用 3 个线程启动
$ pathway spawn --threads 3 python main.py

spawn 命令的入口在 python/pathway/cli.py,同文件还实现了 spawn_from_env(从环境变量读取配置启动)等子命令;新项目的脚手架可参考 README 提及的官方 cookiecutter 模板(独立仓库)。

6.2 Docker 部署(三种姿势)

方式一:基于官方 Pathway 镜像

FROM pathwaycom/pathway:latest

WORKDIR /app

COPY requirements.txt ./
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

CMD [ "python", "./your-script.py" ]

构建并运行:

docker build -t my-pathway-app .
docker run -it --rm --name my-pathway-app my-pathway-app

方式二:单文件脚本直接跑——省去写 Dockerfile:

docker run -it --rm --name my-pathway-app -v "$PWD":/app pathwaycom/pathway:latest python my-pathway-app.py

方式三:标准 Python 镜像 + pip 安装

FROM --platform=linux/x86_64 python:3.10

RUN pip install -U pathway
COPY ./pathway-script.py pathway-script.py

CMD ["python", "-u", "pathway-script.py"]

-u(unbuffered)保证日志实时刷出,容器化部署时推荐保留。

6.3 Kubernetes 与云

README 指出 Docker 容器天然适合在云上以 Kubernetes 部署;面向需要分布式计算与分布式 K8s 部署(含外部持久化)的场景,官方提供 Pathway for Enterprise 版本,定位是端到端数据加工与实时智能分析。此外官方文档提供了几键部署到 Render 等 PaaS 的指引(见 README 外链)。

七、性能定位与基准测试

README 的 Performance 一节声称 Pathway 针对流式与批式处理任务设计为超越 Flink、Spark、Kafka Streaming 等同类技术,并支持其他流框架不易支持的算法(时间连接 temporal joins、迭代图算法、机器学习例程等);官方基准测试位于独立的 pathway-benchmarks 仓库(README 外链),其中包含 WordCount 等对比图。本文不转述具体数值,读者应以该基准仓库的最新公开结果为准。

八、许可证与贡献

  • 许可证:Pathway 采用 BSL 1.1 分发(见 LICENSE.txt),允许无限制的非商业用途及多数免费商业用途;仓库中的代码在 4 年后自动转为 Apache 2.0 开源。Cargo.tomllicense = "BUSL-1.1" 与此一致。部分互补的公共仓库(示例、库、连接器等)采用 MIT 许可。
  • 贡献指引(README Contribution guidelines):若开发了希望集成的库或连接器,建议先以 MIT / Apache 2.0 许可单独发布仓库;核心功能问题鼓励直接提 Issue,并参与官方 Discord 社区讨论。仓库的 CONTRIBUTING.md 有更详细的贡献流程说明。

九、延伸阅读:仓库内值得继续深入的路径

路径 内容
python/pathway/internals/run.py pw.run / pw.run_all 的完整参数与 GraphRunner 启动流程
python/pathway/io 40+ 官方连接器的 Python 实现
python/pathway/xpacks/llm LLM/RAG 工具链(封装、切分、embedding、向量索引)
python/pathway/persistence 状态持久化配置入口
src/engine 差分数据流引擎:dataflow 图、reduce、frontier 时间推进、监控
src/connectors 连接器 Rust 端实现(数据格式解析、存储接入、背压 backlog)
external/differential-dataflow / external/timely-dataflow 增量计算内核的两个核心 crate 本地副本
integration_tests 按连接器与功能组织的集成测试
examples/notebooks / examples/projects 45 个 Notebook 与完整项目级示例
docs/2.developers 仓库内置的开发者文档(159 篇 Markdown)

综上,Pathway 的使用心智模型可以概括为:pw.Schema 描述数据、用 pw.io 连接数据、用表达式与 UDF 声明变换、用一行 pw.run() 交给 Rust 差分引擎增量执行——批与流共用同一套代码,状态、时间与一致性由框架接管。

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

项目优选

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