首页
/ Pathway Live Data Framework 入门全景:从 Schema、连接器到 Rust 引擎的实时数据管道

Pathway Live Data Framework 入门全景:从 Schema、连接器到 Rust 引擎的实时数据管道

2026-09-03 16:18:58作者:郜逊炳

本文基于 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 命名空间一次性暴露了框架的全部核心构件:SchemaTableLiveTableapplysqlgroupbyjoinrunthis(列引用语法)、debugioudfsdemo 等。也就是说,一篇典型的 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 引擎可以基于已知类型做列式存储与向量化计算。

支持的类型包括:

  • 基础类型(详见 数据类型文档):boolstrbytesintfloat
  • 更复杂的类型:Optional 可选类型、时间类型(如 datetime.datetime,内部对应 DateTimeNaive/DateTimeUtc/Duration 等,可从 python/pathway/init.py 的导出中确认)。

Schema 的底层实现位于 python/pathway/internals/schema.py。该文件末尾的 Schema 类通过元类 SchemaMetaclasspython/pathway/internals/schema.py)在子类被定义时收集列声明,__init_subclass__ 还支持 append_onlyid_dtype 等元参数来控制表的主键与追加语义——这些是理解后文"表是动态内容 + 静态结构"模型的关键:主键(id)决定了一行数据在流中是"新增/更新/删除"中的哪一种。

3. Tables:内容动态、结构静态的数据容器

Table 是 Pathway 中真正承载数据的对象:它由若干列组成,每列存放同类型数据,组织方式与关系型数据库的表相似。与静态数据库表不同,Pathway 的表是数据流的快照——Schema 固定不变,内容则随新事件到达而实时更新。核心概念文档中对此有完整论述(见 Core Concepts),其中对新增、更新、删除三类事件如何反映为表行变化的解释值得细读。

在源码层面,表类型定义在 python/pathway/internals/table.py,例如 Table.filterpython/pathway/internals/table.py#L497)与 Table.groupbypython/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/ 目录:csvkafkasqlitegdrivedebeziumpostgreskinesismqttelasticsearchicebergdeltalakes3airbyte 等数十个包,规模远超概览表格所列。文档还特别提到 Airbyte 连接器 可借助 Airbyte 生态连接 300+ 种数据源;完整的可用连接器清单见 连接器总览页

深入看 CSV 连接器 python/pathway/io/csv/init.pyread 签名,可以了解文档没有展开的实用参数:

  • mode"streaming"(默认)或 "static"。streaming 模式下引擎持续监听目录中文件的增、删、改,删除文件会把对应行从表中移除;static 模式则一次性摄取现有数据(两种模式的完整差异见 流式与静态模式文档);
  • schema:结果表的 Schema,也可传 None 让引擎自动推断;
  • csv_settings:CSV 解析器设置(分隔符等);
  • autocommit_duration_ms:默认 1500 毫秒,控制流式模式下自动提交更新的节奏,是延迟与吞吐之间的直接调优旋钮;
  • object_pattern:目录内文件过滤模式;
  • with_metadata:为每行附加 _metadata JSON 列,包含文件的 created_atmodified_atseen_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 中导出的 joinjoin_innerjoin_leftjoin_rightjoin_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 构建在核心算子之上的"标准库"实现,可以精确定位到每个函数:

时态操作的**行为(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/ 下的各子包(kafkapostgrespubsubelasticsearchclickhouseslack 等)。从源码结构看,每个连接器包通常同时包含 readwrite 两侧能力(以支持数据库类系统的读写双向同步),底层共用 internals/datasource.pyinternals/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 工具包与完整学习路径

小结:一条管道的完整骨架

把本文各节拼起来,一个最小但完整的 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 代码声明一条会随数据流动而持续更新的实时管道。

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

项目优选

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