Pathway Live Data Framework 入门全景:从 Schema、连接器到 Rust 引擎的实时数据管道
本文基于 Pathway 官方文档 Live Data Framework 概览 展开,系统介绍如何用 Pathway(pip install pathway 即可安装的 Python 流处理框架)搭建一条完整的实时数据管道:定义数据 Schema、用连接器接入 CSV/Kafka/SQLite 等数据源、执行过滤与聚合等转换、编写时间窗口与时态 Join 等时序算子,最后配置输出连接器并通过 pw.run() 让 Rust 引擎持续运行。读完后你将掌握 Pathway "静态 Schema + 动态表内容 + 增量计算" 的核心编程模型,并能写出可运行、可扩展的流式处理应用。
1. 安装与导入:一条命令接入框架
Pathway 的编程层是纯 Python 的,安装非常简单:
pip install pathway
安装完成后,像导入任意 Python 库一样导入它:
import pathway as pw
从 python/pathway/init.py 的导出列表可以看到,pw 命名空间一次性暴露了框架的全部核心构件:Schema、Table、LiveTable、apply、sql、groupby、join、run、this(列引用语法)、debug、io、udfs、demo 等。也就是说,一篇典型的 Pathway 应用代码只需要 import pathway as pw 这一行导入,就能完成"数据接入 → 转换 → 输出 → 运行"的全部工作。框架的计算引擎则以 Rust 二进制形式随包分发(见 python/pathway/_engine_finder.py 负责定位引擎),这正是文档反复强调"转换在底层用 Rust 编写,因此非常高效"的原因。
2. 定义数据 Schema:静态结构保证类型安全
Pathway 中所有数据表都有 Schema,用于声明列名与列类型,保证数据组织有序、类型一致:
class InputSchema(pw.Schema):
colA: int
colB: float
colC: str
这个 InputSchema 声明了三列:colA(整数)、colB(浮点数)、colC(字符串)。Schema 的作用有两层:一是类型安全,转换操作在运行前就能按声明类型校验;二是运行时优化,Rust 引擎可以基于已知类型做列式存储与向量化计算。
支持的类型包括:
- 基础类型(详见 数据类型文档):
bool、str、bytes、int、float; - 更复杂的类型:
Optional可选类型、时间类型(如datetime.datetime,内部对应DateTimeNaive/DateTimeUtc/Duration等,可从 python/pathway/init.py 的导出中确认)。
Schema 的底层实现位于 python/pathway/internals/schema.py。该文件末尾的 Schema 类通过元类 SchemaMetaclass(python/pathway/internals/schema.py)在子类被定义时收集列声明,__init_subclass__ 还支持 append_only、id_dtype 等元参数来控制表的主键与追加语义——这些是理解后文"表是动态内容 + 静态结构"模型的关键:主键(id)决定了一行数据在流中是"新增/更新/删除"中的哪一种。
3. Tables:内容动态、结构静态的数据容器
Table 是 Pathway 中真正承载数据的对象:它由若干列组成,每列存放同类型数据,组织方式与关系型数据库的表相似。与静态数据库表不同,Pathway 的表是数据流的快照——Schema 固定不变,内容则随新事件到达而实时更新。核心概念文档中对此有完整论述(见 Core Concepts),其中对新增、更新、删除三类事件如何反映为表行变化的解释值得细读。
在源码层面,表类型定义在 python/pathway/internals/table.py,例如 Table.filter(python/pathway/internals/table.py#L497)与 Table.groupby(python/pathway/internals/table.py#L1192)都是该类的实例方法。文档概览示例中的 pw.this.colA 语法,就是 python/pathway/internals/column_namespace.py 中列引用机制的体现,让表达式写法接近 SQL 而不是裸函数调用。
4. 输入连接器:把数据源接成实时表
要创建表,需要连接器(connector)。连接器负责实时从外部数据源读取并摄取数据。以 CSV 为例:
input_table = pw.io.csv.read('./data/', schema=InputSchema)
文档概览中给出的常用输入连接器一览:
| 输入连接器 | 示例 |
|---|---|
| CSV | pw.io.csv.read('./data/', schema=InputSchema) |
| Kafka | pw.io.kafka.read(rdkafka_settings, topic="example", schema=InputSchema, format="csv") |
| SQLite | pw.io.sqlite.read('./data_path/', table_name, schema=InputSchema) |
| Google Drive | pw.io.gdrive.read(object_id='***', service_user_credentials_file="credentials.json") |
这些连接器在源码中都有对应的独立子包,位于 python/pathway/io/ 目录:csv、kafka、sqlite、gdrive、debezium、postgres、kinesis、mqtt、elasticsearch、iceberg、deltalake、s3、airbyte 等数十个包,规模远超概览表格所列。文档还特别提到 Airbyte 连接器 可借助 Airbyte 生态连接 300+ 种数据源;完整的可用连接器清单见 连接器总览页。
深入看 CSV 连接器 python/pathway/io/csv/init.py 的 read 签名,可以了解文档没有展开的实用参数:
mode:"streaming"(默认)或"static"。streaming 模式下引擎持续监听目录中文件的增、删、改,删除文件会把对应行从表中移除;static 模式则一次性摄取现有数据(两种模式的完整差异见 流式与静态模式文档);schema:结果表的 Schema,也可传None让引擎自动推断;csv_settings:CSV 解析器设置(分隔符等);autocommit_duration_ms:默认 1500 毫秒,控制流式模式下自动提交更新的节奏,是延迟与吞吐之间的直接调优旋钮;object_pattern:目录内文件过滤模式;with_metadata:为每行附加_metadataJSON 列,包含文件的created_at、modified_at、seen_at等时间戳。
5. 转换(Transformations):用 Python 描述增量计算
数据接入后,就可以用 Pathway 的转换操作定义处理逻辑。官方示例演示了"过滤 + 分组求和":
filtered_table = input_table.filter(input_table.colA > 0)
result_table = (
filtered_table
.groupby(filtered_table.colB)
.reduce(sum_val=pw.Reducers.sum(pw.this.colC))
)
这段代码先保留 colA > 0 的行,再按 colB 分组,对每组的 colC 求和。值得注意的是语义:在流式场景下,这不是一次性计算,而是增量维护的聚合——后续任何新到达或更新的数据到达时,Rust 引擎只重算受影响的分组。
框架支持的操作面较广,概览文档归纳如下:
| 类别 | 操作 | 示例 |
|---|---|---|
| 算术运算 | +、-、*、/、//、%、** |
t.select(new_col = t.colA + t.colB) |
| 比较运算 | ==、!=、<、<=、>、>= |
t.select(new_col = t.colA <= t.colB) |
| 布尔运算 | &(AND)、|(OR)、~(NOT)、^(XOR) |
t.select(new_col = t.colA & (t.colB < 3)) |
| 过滤 | filter |
t.filter(pw.this.column > value) |
| 应用函数 | pw.apply |
t.select(new_col=pw.apply(func, pw.this.colA)) |
| SQL 命令 | pw.sql |
pw.sql(query, tab=t) |
其中 pw.apply 允许嵌入任意 Python 函数(即 UDF),pw.sql 则提供 SQL 接口(见 SQL 文档)。基本操作全集见 表操作指南。
框架还提供更高级的转换原语:
- Group-by 与聚合:见 Group-by / Reduce 手册,内置
pw.Reducers提供 sum、min、max、mean 等聚合器; - Join 操作:见 Join 手册,支持 inner/left/right/outer 等连接模式(
python/pathway/__init__.py中导出的join、join_inner、join_left、join_right、join_outer即为对应入口)。
6. 时态转换:窗口、ASOF Join 与 Interval Join
作为数据流处理框架,Pathway 内置一组时态(temporal)操作,概览文档列出了四类:
| 类别 | 操作 | 示例 |
|---|---|---|
| 窗口操作 | windowby(sliding / tumbling / session 窗口) |
t.windowby(t.time, window=pw.temporal.tumbling(duration=...), ...).reduce(...) |
| ASOF now join | asof_now_join |
t1.asof_now_join(t2, t1.t, t2.t, t1.name == t2.name, how=..., direction=...).select(...) |
| Interval join | interval_join(outer/left/right) |
t1.interval_join(t2, t1.t, t2.t, pw.temporal.interval(...), t1.col == t2.col).select(...) |
| Window join | window_join(outer/left/right) |
t1.window_join(t2, t1.t, t2.t, pw.temporal.sliding(...), t1.col == t2.col).select(...) |
这些算子在源码中位于 python/pathway/stdlib/temporal/ 目录,是纯 Python 构建在核心算子之上的"标准库"实现,可以精确定位到每个函数:
windowby、sliding、tumbling窗口定义:python/pathway/stdlib/temporal/_window.py;interval_join及其 inner/left/right/outer 变体:python/pathway/stdlib/temporal/_interval_join.py;asof_now_join及其变体:python/pathway/stdlib/temporal/_asof_now_join.py。
时态操作的**行为(behavior)**决定了"准确性、延迟、内存消耗"三者之间的权衡,可以用 asof(事件到达时才处理)与 eager(数据一到达就提前计算)两种策略及其组合来调优,详见 时态行为文档 与 窗口行为配置,窗口操作本身参见 窗口手册、Interval Join 与 Window Join。ASOF now join 在 索引机制文档 中也有讲解。
7. 配置输出:把结果送回外部系统
管道处理完成后,用输出连接器把结果写出框架。最简单的是写 CSV 文件:
pw.io.csv.write(result_table, './output/')
概览文档列出的常用输出连接器:
| 输出连接器 | 示例 |
|---|---|
| CSV | pw.io.csv.write(table, './output/') |
| Kafka | pw.io.kafka.write(table, rdkafka_settings, topic_name="example", format="json") |
| PostgreSQL | pw.io.postgres.write(table, output_postgres_settings, "sum_table") |
| Google PubSub | pw.io.pubsub.write(table, publisher, project_id, topic_id) |
与输入侧一样,输出实现同样对应 python/pathway/io/ 下的各子包(kafka、postgres、pubsub、elasticsearch、clickhouse、slack 等)。从源码结构看,每个连接器包通常同时包含 read 与 write 两侧能力(以支持数据库类系统的读写双向同步),底层共用 internals/datasource.py 与 internals/datasink.py 中的源/汇抽象。完整清单仍建议查阅 连接器总览页。
8. 运行管道:pw.run() 与"永久运行"语义
一切就绪后,一行命令启动计算:
pw.run()
这里有一个必须建立的认知:pw.run() 启动的进程会一直监听数据源的新更新,直到进程被终止——计算永远运行下去才是框架的正常行为,而不是处理完初始数据就退出(这一设计动机在 Core Concepts 的"用 Rust 引擎运行计算"一节有解释)。从源码看,run 函数由 python/pathway/internals/api.py 提供并经 python/pathway/init.py 导出,其职责是把前面声明式的管道(连接器 + 转换构成的图)交给 Rust 引擎执行。这也解释了 Pathway 的核心架构特征:管道定义与计算执行严格分离——你在代码里声明的是"处理配方"(数据源、转换、输出位置),真正的数据流动发生在 pw.run() 之后,由 Rust 引擎按增量语义持续推进。
9. 延伸方向:LLM 工具包与完整学习路径
- LLM xpack:Pathway 提供 LLM 扩展包(源码位于 python/pathway/xpacks/),可在框架内构建嵌入、检索、问答等 LLM 管道,详见 LLM xpack 概览;
- 入门实操:从 第一个实时应用 动手写端到端示例;
- 深入概念:Core Concepts 完整覆盖连接器、表、转换、输出、数据流与运行时六大核心概念;
- 框架对比:为什么选择 Live Data Framework 与 批处理模式 可帮助理解它与批式 ETL 的差异。
小结:一条管道的完整骨架
把本文各节拼起来,一个最小但完整的 Pathway 实时应用就是五步:
import pathway as pw
# 1. 定义 Schema
class InputSchema(pw.Schema):
colA: int
colB: float
colC: str
# 2. 输入连接器(流式监听 ./data/ 目录)
input_table = pw.io.csv.read('./data/', schema=InputSchema)
# 3. 转换:过滤 + 分组聚合
filtered_table = input_table.filter(input_table.colA > 0)
result_table = (
filtered_table
.groupby(filtered_table.colB)
.reduce(sum_val=pw.Reducers.sum(pw.this.colC))
)
# 4. 输出连接器
pw.io.csv.write(result_table, './output/')
# 5. 运行:持续监听数据源,直到进程终止
pw.run()
这套"Schema 声明 → 连接器建表 → 转换描述计算 → 输出连接器 → pw.run()"的固定骨架,正是 Pathway 文档 Live Data Framework 概览 所传递的核心编程范式:用静态 Python 代码声明一条会随数据流动而持续更新的实时管道。
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