Pathway Live Data Framework:基于 Rust 差分计算引擎的 Python 实时流批一体 ETL 框架
本文以仓库根目录的 README.md 为骨架,系统讲解 Pathway(Pathway Live Data Framework)的定位、核心概念与上手路径,并结合仓库中的 pyproject.toml、Cargo.toml 与 python/pathway 源码,深入剖析 pw.run() 运行参数、连接器体系、LLM 扩展、持久化机制以及本地/Docker/Kubernetes 三种部署方式,帮助读者掌握这套“Python 写逻辑、Rust 跑引擎”的流式数据框架的完整使用方法。
一、什么是 Pathway:流批一体的 Python ETL 框架
README.md 对 Pathway 的定义是:一个面向流处理(stream processing)、实时分析(real-time analytics)、LLM 管线和 RAG 的 Python ETL 框架。它的设计目标可以归纳为三点:
- 易用:提供易用的 Python API,可以无缝集成常用的 Python ML 库;
- 同一份代码跨越开发与生产:无论是本地开发、CI/CD 测试、批处理作业、流回放(stream replay)还是实时流处理,都可以使用同一段管线代码;
- 性能由 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.toml 中
differential-dataflow = { path = "./external/differential-dataflow" }与 Cargo.toml 中timely = { path = "./external/timely-dataflow/timely" },即仓库在 external/differential-dataflow 和 external/timely-dataflow 内置了这两个核心 crate 的本地副本; - Python 与 Rust 之间通过 PyO3 桥接,Cargo.toml 显示
pyo3启用了abi3-py310特性,与 Python 侧最低 3.10 的版本要求一致;打包由 maturin 完成(见 pyproject.toml 与 pyproject.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.toml 中
requires-python = ">=3.10"完全一致; - 平台支持 MacOS 与 Linux,其他系统需通过虚拟机运行(README 原文带 ⚠️ 提示)。
安装命令:
pip install -U pathway
除了核心依赖,pyproject.toml 还声明了 pandas、numpy、pyarrow、fastapi、uvicorn、OpenTelemetry 等基础依赖,并提供了按功能拆分的可选依赖组(extras),常用的几组包括:
| 安装组 | 用途 | 关键依赖(摘自源码) |
|---|---|---|
pathway[sql] |
pw.sql 字符串 SQL 表达式 |
sqlglot==10.6.1(新版与当前 pw.sql 实现不兼容,故精确锁定) |
pathway[xpack-llm] |
LLM 管线扩展 | openai、litellm、langchain、llama-index-core 等 |
pathway[xpack-llm-local] |
本地推理 | sentence_transformers、transformers |
pathway[xpack-llm-docs] |
文档解析 | docling、unstructured、paddleocr 等 |
pathway[milvus] |
Milvus 向量库连接器 | pymilvus、milvus-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()
各步骤对应的仓库实现:
- Schema 定义:
pw.Schema定义在 python/pathway/schema.py。声明 Schema 可以让引擎提前知道列类型,走更快的解析路径;不声明时引擎会做类型推断。 - CSV 连接器:位于 python/pathway/io/csv,配套的 Rust 端数据格式解析在 src/connectors/data_format(
dsv.rs、csv等)。 filter/reduce:reduce依赖 python/pathway/reducers.py,其中sum、min、max、count、avg、count_distinct_approximate等聚合器在 python/pathway/internals/reducers.py 中实现(如 python/pathway/internals/reducers.py 的sum)。这些聚合器以可加/可逆的增量式实现,这正是差分数据流能做“只重算变化部分”的基础。pw.run():见下一节完整参数解析。- 结果写出:
pw.io.jsonlines.write位于 python/pathway/io/jsonlines。
除 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.py,run 的完整签名与参数如下(全部为关键字参数):
| 参数 | 默认值 | 说明 |
|---|---|---|
debug |
False |
是否启用 table.debug() 算子的输出 |
monitoring_level |
MonitoringLevel.AUTO |
监控粒度:NONE / IN_OUT / ALL;AUTO 时框架根据输出端是否交互自动在 NONE 与 IN_OUT 间选择 |
with_http_server |
False |
是否启动带运行时指标的 HTTP 服务(源码注释标注“即将弃用”,监控已演进为内建 dashboard) |
default_logging |
True |
是否允许 Pathway 自行设置 logging handler;想接管日志时置为 False |
persistence_config |
None |
持久化配置,用于在更新或崩溃后重启管线时恢复计算状态(对应 python/pathway/persistence 与 src/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_all(python/pathway/internals/run.py),功能相同但禁用 tree-shaking 优化,即不做“无用输出剪枝”,适合需要完整执行所有注册输出的调试场景。
五、核心特性逐项解析
README.md 的 Features 一节列出了六大特性,这里结合仓库证据逐项展开:
5.1 丰富的连接器生态
README 声明内置 Kafka、GDrive、PostgreSQL、SharePoint 等连接器,且 Airbyte 连接器可接入 300+ 数据源;缺失的连接器可用 Python 自定义。
从源码结构看,python/pathway/io 下的连接器目录规模远超 README 的举例,按字母序包括:airbyte、bigquery、chroma、clickhouse、csv、debezium、deltalake、duckdb、dynamodb、elasticsearch、gdrive、iceberg、jsonlines、kafka、kinesis、leann、logstash、milvus、minio、mongodb、mqtt、mssql、mysql、nats、pinecone、plaintext、postgres、pubsub、pulsar、qdrant、questdb、rabbitmq、redpanda、s3、slack、sqlite、weaviate 等 40 余个子目录。Rust 端 src/connectors/data_storage 则实现了 AWS(S3)、数据湖(Delta Lake/Iceberg)等存储接入;Cargo.toml 的依赖列表也印证了这些底层客户端(rdkafka、pg_walstream、iceberg、elasticsearch、questdb-rs、rumqttc、pulsar 等)都被直接编进引擎。仓库还在 integration_tests 下按连接器组织了对应集成测试(kafka/、db_connectors/、s3/、airbyte/ 等),可用来验证各连接器的真实读写行为。
5.2 无状态与有状态变换
Pathway 支持 join、窗口(windowing)、排序等有状态变换,且很多变换直接在 Rust 中实现以提升性能;同时任何 Python 函数都可以作为 UDF 参与计算——可以自写,也可以直接调用任意 Python 库。从 src/engine 的模块划分(dataflow/、reduce.rs、expression.rs 等)可以看出,聚合、join 等状态算子确实由引擎原生实现,而 Python UDF 则通过 pyo3 回调执行。
5.3 持久化(Persistence)
README 强调持久化可保存计算状态,使管线在更新或崩溃后能够重启恢复。实现上,Python 侧 pw.run(persistence_config=...) 接收 python/pathway/persistence 中的 Config;Rust 侧 src/persistence 包含 state.rs、operator_snapshot.rs、frontier.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.toml 的
xpack-llm/xpack-llm-local依赖组(openai、litellm、cohere、langchain、llama-index-core、sentence_transformers等); - Rust 端 src/external_integration 实现了内存向量检索(
brute_force_knn_integration.rs、usearch_integration.rs)与全文检索(tantivy_integration.rs、qdrant_integration.rs),Cargo.toml 中的usearch、tantivy依赖即服务于此; - 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.toml 的
license = "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 差分引擎增量执行——批与流共用同一套代码,状态、时间与一致性由框架接管。
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 StartedRust0622
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