首页
/ Pathway 批量处理(Batch Processing)实战指南:同一套 Rust 引擎与 Python API 平滑打通静态与流式工作负载

Pathway 批量处理(Batch Processing)实战指南:同一套 Rust 引擎与 Python API 平滑打通静态与流式工作负载

2026-09-07 14:52:15作者:龚格成

导读

Pathway Live Data Framework 的核心是一个统一的数据流引擎:面对静态数据,它以 static mode(静态模式) 将数据作为一个批次一次性处理完毕(即传统意义上的 batch processing);面对持续到达的数据,它则进入 streaming mode(流式模式) 做增量计算。本文以官方开发者文档 Batch Processing in Python 为主线,从"为什么要用 Pathway 做批量处理"出发,结合仓库内概念文档、连接器文档、持久化文档与真实示例代码,帮助你掌握:静态模式与流式模式的切换方式、Rust 引擎带来的性能与可扩展性、如何用持久化支撑 backfilling 与周期性计算,以及如何用最小的代码改动从批处理平滑迁移到实时处理。


什么是 Pathway 中的 Batch Processing:统一引擎的两种模式

Pathway 的引擎用 Rust 实现,并同时驱动静态与实时两类数据负载。从 Streaming and Static Modes 相关概念 可以看出框架对两种模式的定义:

  • Static mode(静态模式,即批处理):一次性加载全部数据并完成处理,处理完毕后进程自然结束,不会像流式那样持续等待新数据。文档原话是 "In the static mode, all the data is loaded and processed at once and then the process terminates"(见 50.concepts.md 的 Static mode 小节)。
  • Streaming mode(流式模式):引擎持续监听数据源,每当新数据到达,就基于当前已计算的版本做增量更新。框架正常的行为是"一直运行,直到进程被终止"。

也就是说,"batch processing" 在 Pathway 中并不是一套独立于流式之外的子系统,而是同一个引擎、同一张计算图(dataflow)在不同数据源模式下的表现。计算图本身是惰性构建的:调用 selectjoinreduce 等操作符时只描述意图,并不真正执行;只有当 pw.run()(或调试用 API)启动后,引擎才开始按数据流推演计算。正因为底层共享同一张计算图与同一套内核,静态模式与流式模式之间几乎没有语义差异,唯一的例外是少数只在流式上下文里才有意义的时间相关操作。

由此可以理解官方推荐的使用路径:

  • 在批处理阶段起步:用静态文件(CSV、Parquet、静态表)把业务逻辑打磨正确,适合每日一次的离线计算或先期的算法验证;
  • 数据量增长后平滑转流:把数据源从静态改为流式(Kafka、实时文件目录、数据库变更流等),计算逻辑一行都不用改;
  • 上线前先在静态数据上测试:例如 从 Jupyter 部署到生产环境 中描述的工作流,先在 Notebook 里用静态模式验证,再平滑过渡到真实环境的流式管道,从而保证批、流两侧行为一致。

一、统一框架:批量与流式共用同一套 Python API

只换连接器,不换业务逻辑

原文档强调的最核心卖点是 "Unified Framework for Stream and Batch Processing":Pathway 为批处理和流处理提供统一的 Python API,切换两种形态的唯一动作是更换数据源连接器,其余管道代码保持不动。仓库中的专门教程 Switching from Batch to Streaming 用一个"读取 CSV → 求和 → 写出"的例子完整演示了这一过程。

先看静态(批处理)版本——读取 CSV 目录、对整批数据求 value 之和、写出结果:

import pathway as pw

# WRITE SOME STATIC CODE

# read data
class InputSchema(pw.Schema):
    value: int

t = pw.io.csv.read(
    './sum_input_data/',
    schema=InputSchema,
    mode="static",
)

# process data
t = t.reduce(pw.reducers.sum(t.value))

# write data
pw.io.csv.write(t, "output.csv")

# run
pw.run()

注意关键参数 mode="static"——正是它让 CSV 读取器进入静态模式,触发一次性整批处理。之后想要实时化,只需把该参数改成 "streaming"

t = pw.io.csv.read(
    './sum_input_data/',
    schema=InputSchema,
    mode="streaming",
)

其余实现完全不变:读、算、写、运行依旧是同一份代码。CSV 文件输入连接器同时支持两种模式,所以它能充当从批到流的最小演示载体。需要说明的是,并非所有连接器都同时支持静态与流式两种模式,动手前建议对照 支持的输入/输出数据源清单 选择合适的数据源。

切换到流式后,实际发生的变化只有一点:输出不再是一份固定的静态结果,而是一个不断被修订的数据流。每当有新 CSV 文件落入输入目录,引擎会摄取新值并自动更新求和结果,而 reduce 这类聚合不会从零重算,只对新数据做增量推进——这正是引擎数据流模型的天然属性(可参考 50.concepts.md 中关于数据流维护"最新版本数据、只更新相关部分"的描述)。

仓库内的真实静态管道示例

这一模式并非停留在文档层面。仓库 examples/projects/option-greeks/greeks-static.py 就是一个以静态 CSV 为输入计算期权 Greeks 的完整项目,其开头正是用 mode="static" 读取定义表:

class DefinitionInputSchema(pw.Schema):
    ts_recv: int  # Time in ns when the data was received
    raw_symbol: str  # symbol of option
    expiration: int  # expiration time of the option
    instrument_class: str  # type of option
    strike_price: float  # see below for what it means
    underlying: str  # symbol of the first underlying instrument
    instrument_id: int  # An identifier of the option

table_es = pw.io.csv.read(definitions_path, schema=DefinitionInputSchema, mode="static")

从示例命名(greeks-static.py)与代码注释 "The static case: the data is read from CSV files" 可以推断,该项目刻意先固化"静态数据读入、本地验证、复用同一 schema 与处理逻辑"的开发路径,之后再演化出相应的流式版本。这也是原文档建议的工程节奏:先批后流、批流共码,避免为了不同处理类型维护多套框架、多份代码,从而降低数据工程复杂度。

保持统一代码库的收益

按原文档的归纳,批流共用一套代码直接带来三方面收益:

  1. 开发与测试效率:先在静态数据上验证管道逻辑的正确性(不必关心时间与一致性语义),通过后再上线实时场景;
  2. 运维简化:单一代码库、单一执行模型,无需在"批处理框架 + 流处理框架"两套系统间来回搬运逻辑;
  3. 结果一致性:除少数仅对流式有意义的时序算子外,引擎保证流式输出等价于"用批处理方式重算已接收数据"所得的结果,减少批流不一致带来的排查成本。

二、可扩展性与性能:Rust 内核、无 GIL、支持高级算子

为什么 Rust 引擎对批处理同样重要

很多人认为批处理是低频任务,对性能不敏感。但原文档指出,Pathway 的引擎直接用 Rust 构建(无 JVM),从根本上绕开了 Python 在原生计算上的瓶颈——包括著名的全局解释器锁(GIL)——从而让单机即可借助**多进程(multi-processing)**扩展到大数据量处理。其机制是:Python 仅用于声明类型化的管道(type the pipeline),一旦计算图就绪,数据表的存储与绝大多数变换都由 Rust 引擎接管(50.concepts.md 中关于引擎的论述与此一致)。

对批处理场景而言,这至少意味着:

  • 处理吞吐与数据集规模不再受 GIL 拖累:即便一次灌入整批静态数据,核心计算仍在内核级并行下完成;
  • 内存效率更高:表内容由 Rust 侧管理,配合数据流只保留当前版本与必要增量,减少重复计算开销;
  • 多机/多进程可扩展路径明确:仓库在 高级开发文档 目录下提供了 worker 架构与多机运行的专题说明,可作为把批/流管道从原型推向规模的延伸阅读。

面向 ML 的迭代计算与时间算子

除了常规关系算子,原文档特别强调 Pathway 提供了两类对批处理(尤其是机器学习与时间型业务)很有价值的高级操作

  • 迭代计算(iterative computations):常见于机器学习中的不动点迭代、PageRank、图算法等需要反复在整张表上收敛计算的场景;
  • 时间算子(temporal operations):虽然部分时间语义只在流式下有意义,但 Pathway 的时间感知能力(如 asof join、interval join 等,参见 temporal-data 相关文档)同样能服务于带时间戳的静态数据集。

这些算子以 Rust 为后端的统一数据流实现,意味着即使是批量迭代,也沿用了同一套增量更新机制,从架构层面避免"框架割裂、能力不对齐"的常见痛点。


三、持久化(Persistence):为 backfilling 与周期性计算而生

增量而非重算:批量场景下的持久化价值

原文档将持久化列为 Pathway 适合批处理的第三大理由,逻辑链条如下:

  1. Rust 引擎天然设计为增量计算——只处理数据的变化量,而不是重算整个数据集;
  2. 静态模式下反复做周期性批计算(例如每日统计、定期回填)时,这一特性直接转化为收益;
  3. 配合持久化保存计算状态,Pathway 每次只处理新增或更新的数据部分,而不是从头开始全量计算,对大型数据集可显著节省计算资源与时间。

这正是"backfilling"的典型语义:历史数据缺失或延迟到达后,需要在不重放全部历史的前提下把空缺补上。仓库 persistence_restart_with_new_data 文档 描述的场景与此高度吻合:实时日志可能按小时增量到达,或由 cron 类任务每十分钟/小时/天周期性采集,每次新数据出现都需重建分析结果——若每次都对全量日志重算,成本将随数据增长线性膨胀。

让管道具备持久化能力的最小改动

持久化官方文档 给出了一个可复制的 wordcount 例子。先看未持久化版本:

import pathway as pw

class InputSchema(pw.Schema):
    word: str

words = pw.io.csv.read("inputs/", schema=InputSchema)
word_counts = words.groupby(words.word).reduce(words.word, count=pw.reducers.count())
pw.io.jsonlines.write(word_counts, "result.jsonlines")
pw.run()

这个程序一旦重启就会从头扫描文件、重新统计。要让它"断点续算",只需两步:

第一步,创建持久化后端与配置(把中间状态存到本地 state 目录):

persistence_backend = pw.persistence.Backend.filesystem("./state/")
persistence_config = pw.persistence.Config(persistence_backend)

第二步,把配置传给 pw.run

pw.run(persistence_config=persistence_config)

改造后程序变为持久化的:若被重启,它会从上次停止的位置继续计算与输出,而不是推倒重来。配合 snapshot_interval_ms 参数(配置希望保持的快照新鲜度,毫秒为单位),还可以在"状态快照保存成本"与"故障恢复粒度"之间做权衡。在 python/pathway/tests/test_persistence.py 中可以看到对应的自动化验证:test_groupby_count(第 49 行起)与 test_persistence_modifications(第 222 行起)都把 pw.persistence.Config 传入 pw.run,从测试层面确认了分组聚合等有状态计算在持久化配置下可正确恢复与更新。

需要补充的关键细节(同样来自持久化文档):

  • 元数据与快照分离pw.persistence.Config 内部区分 metadata storage(记录已推进的时间、各数据源当前读取位置、计算图描述,体量不随数据量增长)与 snapshot storage(数据快照,体量随数据量增长);二者都可通过 pw.persistence.Backend 配置为本地文件系统或 S3。
  • 数据源唯一标识(name):持久化恢复依赖给输入源指定稳定的 name。可以自动生成,也可以手动指定,例如 pw.io.csv.read("inputs/", schema=InputSchema, name="words_source")。只要 name 不变,即使数据源路径变化、格式变化(如 JSON 换 CSV)或字段改名,引擎仍能识别这是同一张表。
  • 一致性语义:当前实现的恢复保证是 at-least-once(至少一次);输入在内部被切分为小事务批次,若执行被中断,未闭合事务批次对应的输出可能重复,而优雅终止场景可以做到 exactly-once。

在周期性批任务与 backfilling 中,这套机制的实际效果是:只计算新增/变化的部分,框架自动记住已经处理到的数据位置与计算结果,从而让"每日统计""每小时回填""增量补数"等任务在数据规模变大后依旧轻量。


四、工作流建议:先批后流,架构稳健可扩展

原文档的结论部分实际上给出一条工程建议:受迫式迁移(forced migration)代价高昂,不要等到实时处理成为刚需时才被动重构。既然 Pathway 允许在静态数据上以批处理方式开发、验证并部署,之后通过替换连接器平滑切换到流式,就应当:

  1. 从一开始就用 Pathway 的统一 API 编写管道,哪怕当前只需要每日一次的批处理;
  2. 在静态模式下完成充分的正确性验证(本地文件、静态 CSV 均可用作输入);
  3. 当数据量增长、业务对时效性提出要求时,再把数据源切到流式(Kafka、实时文件目录等),让同一套代码继续服役。

这种"批流一体"的架构既规避了双框架的运维负担,也让管道在吞吐与数据规模增长时保持鲁棒。进一步迁移的逐行代码示例与踩坑说明,可直接参阅 Switching from Batch to Streaming 与部署章节中 持久化恢复教程 等专题。


总结

回到原文档的核心命题:Pathway Live Data Framework 不仅在实时流处理上表现突出,其架构特性同样使它成为静态数据批处理的强力选项

  • 统一引擎与统一 API:批处理 = 同一引擎的 static mode,切换批/流只需更换数据源连接器,管道主体代码不动,天然支持"先批后流"的渐进式演进;
  • Rust 内核的性能与扩展性:无 GIL、可多进程扩展,并原生支持迭代计算与时间类算子,适配 ML 与时间型批量负载;
  • 持久化支撑 backfilling 与周期计算:增量计算 + 状态持久化让批量重算变成"只处理增量",为大型数据集的每日统计、历史回填显著节省资源;
  • 平滑迁移路径:从静态起步、充分测试,待实时需求出现后再切换到流式,可避免被迫重构带来的工程风险。

如果你正在为一个可能走向实时的业务设计管道,用 Pathway 的静态模式做批处理,就是为未来的流式化提前铺好同一条轨道。相关一手资料可按如下路径继续深挖:

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