首页
/ Pathway:把批处理流水线一键切换为流处理,只差一个 mode 参数

Pathway:把批处理流水线一键切换为流处理,只差一个 mode 参数

2026-09-05 18:30:51作者:滑思眉Philip

在 Pathway(Pathway Live Data Framework)中,"从批处理切换到流处理"不是一个架构重构问题,而是一次连接器参数的变更。本文围绕用户指南中的 切换批到流 一文展开:先用一个完整的 CSV 求和示例演示静态(static)批处理写法,再仅改动数据源的 mode 参数把它变为实时流处理,并结合 CSV 连接器源码 深入解析两种模式在引擎层面的差异、输出从静态结果变为数据流后应如何理解,以及用 pw.demo.range_stream 生成人工数据流进行本地验证的方法。读完本文,你可以掌握在 Pathway 中以最小改动把已验证的批处理流水线升级为实时流水线的完整方案。

核心思路:只换连接器,逻辑零改动

Pathway 是一个统一的批处理和流处理框架(unified batch and streaming processing framework)。官方文档给出的结论非常直接:你的静态流水线已经写好、测试并验证过,现在要让它跑起来实时数据,唯一需要做的事情就是把数据源连接器换成对应的流式连接器,其余交给 Pathway 处理。

# 静态批处理 -> 流处理,唯一改动
t = pw.io.csv.read("./sum_input_data/", schema=InputSchema, mode="streaming")

CSV 连接器实现 看,read 的函数签名中 mode 参数被显式声明为 Literal["streaming", "static"],且默认值就是 "streaming"——也就是说在 Pathway 中"流"是一等公民,批反而是静态特例:

def read(
    path: str | PathLike,
    *,
    schema: type[pw.Schema] | None = None,
    csv_settings: CsvParserSettings | None = None,
    mode: Literal["streaming", "static"] = "streaming",  # 默认即流模式
    ...
) -> Table:

同时要注意:并非所有连接器都同时支持 static 和 streaming 两种模式。选型时先查 支持的数据源列表,用法细节可参考 Input/Output API 文档;Pathway 的输入输出实现统一位于 python/pathway/io 目录,涵盖 CSV、文件、数据库、消息队列、向量库、LLM 组件等数十种连接器,每个连接器的 read 文档字符串都标注了它在两种模式下是否可用。

第一步:写一段静态批处理代码

从静态数据起步是最省事的路径:开发、测试时不用操心时间和一致性问题,可以先确保管道逻辑本身正确。下面是一个完整可运行的最小示例——读取 CSV 文件、对所有值求和、把结果写回新的 CSV 文件:

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()

只需要准备一个目录 ./sum_input_data/,里面放若干带单列 value 的 CSV 文件即可。运行 pw.run() 后,output.csv 中就是全部数值之和。

这段代码的三个要素在切流之后一个都不用改

  • InputSchema:用 Python 类声明表结构(value: int),是读写两侧共享的契约;
  • t.reduce(pw.reducers.sum(t.value)):声明式的聚合逻辑,不关心数据是一次到达还是持续到达;
  • pw.io.csv.write(t, "output.csv") + pw.run():声明输出并启动引擎。

第二步:把 mode 改为 streaming

现在你希望"往目录里不断新增 CSV 文件,总和自动实时更新"——即从静态数据源切换为流式数据源。在 Pathway 中,这个改动就是一行:

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

完整代码对比如下(两处 mode 取值是唯一差异):

import pathway as pw

# WRITE SOME STREAMING CODE

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

t = pw.io.csv.read(
    './sum_input_data/',
    schema=InputSchema,
    mode="streaming",   # 原 "static" 改为 "streaming"
)

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

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

#run
pw.run()

"就这样,没了。" 剩下的实现保持不变。

mode 参数在引擎层做了什么

结合 CSV 连接器源码mode 参数的文档字符串,两种模式的语义差异是明确的:

  • mode="static":引擎只考虑当前已有的数据,把全部数据一次性摄入一个 commit,处理完即结束。适合离线开发、测试与校验。
  • mode="streaming":引擎会持续等待目录中的更新,追踪文件的新增、删除、修改并把这些事件反映到表状态中。例如某个文件被删除后,该文件读出的行也会被从表中移除——这意味着聚合结果会正确地"减去"被删数据的贡献,而不是残留脏值。

从源码结构看,pw.io.csv.read 只是对 文件系统连接器 的薄封装(转发 format="csv"mode 等参数),所以"批流一键切换"的能力实际上来自底层 fs 连接器对目录轮询/事件追踪的统一实现,CSV 只是其中一种解析格式。

顺带掌握的 CSV 读取参数

既然改动就发生在 pw.io.csv.read 上,把其余参数一并了解清楚会更有实战价值(默认值均取自 源码签名):

参数 默认值 说明
path 必填 文件、目录或 glob 模式。传入目录时,文件按修改时间升序处理
schema None 结果表的 schema 类
mode "streaming" "streaming"(持续追踪变更)或 "static"(一次性摄入)
csv_settings None CSV 解析器设置(分隔符等)
object_pattern "*" 目录内文件的 Unix shell 风格过滤(官方建议改用 path 中的 glob,该值将弃用)
with_metadata False 为 true 时额外加一列 _metadata(JSON),含 created_at / modified_at / seen_at / owner / path / size
autocommit_duration_ms 1500 两次 commit 之间的最大间隔。每隔该毫秒数,连接器收到的更新就被提交并推入 Pathway 的计算图
max_backlog_size None 限制任一时刻从源读取并保留在处理中的条目数。达到上限后暂停读取,直到部分条目处理完成。适合初始数据突刺很大的源,避免内存峰值
name None 连接器唯一名称,用于日志与监控面板;启用持久化时同时作为快照名
debug_data 调试模式下替代原始数据的静态数据

对流式场景尤其值得注意的两个参数:autocommit_duration_ms 决定了"新 CSV 文件落盘后多久内被计入结果"(默认 1.5 秒一次批量提交,而非逐条即时),max_backlog_size 则帮助你在面对初始海量文件时控制内存水位。

切换之后:输出从静态结果变成了数据流

恭喜,你的静态项目已经是实时数据处理管道了。但对使用者来说,实际变化有多大?答案是:不大,Pathway 替你处理了一切——你不需要自己管理数据的时序性,迟到(late)和乱序(out-of-order)的数据点由框架管理。

需要真正理解的只有这一点:流式系统的输入是永无止境的数据流,因此 Pathway 先基于当前可用数据算出一份输出,此后每有新数据到达就修订(revise)结果。输出不再是静态结果,而是一条数据流。

这一点在输出连接器上体现得很直观。csv.write 的文档示例 展示了 output.csv 的真实形态:

age,owner,pet,time,diff
10,"Alice","dog",0,1
9,"Bob","cat",0,1
8,"Alice","cat",0,1

除了数据列外,还多出两个列:

  • time:该行是在第几个操作 minibatch 被读入/更新的。静态数据全部在初始批次(time=0)进入;流式运行时,每新增一个 CSV 文件,修订过的聚合值会以递增的 time 再次出现;
  • diff:该行在该批次的变更量。初始读入为 1;后续修订会以 ±1 等差值形式表达"删旧行、加新行"的更新语义。

这就是批流切换后唯一的真实变化:pw.io.csv.write 写出的是更新流,而不是最终结果文件。每一轮新 CSV 文件被摄入后,reduce(sum) 的结果被重新计算,output.csv 中追加一组带新 time 的记录。更多输出语义可参考 流式与静态模式 与 第一个实时应用 中对输出结果的解释。

mode="static" 的完整语义(一次 commit、结果即终态),官方单独有一篇 静态模式/批处理 说明,可与本文的流式切换对照阅读。

本地验证:用人工数据流替代文件目录

如果你暂时不想手动往目录里丢文件,Pathway 提供了生成合成流数据的工具。pw.demo.range_stream 会生成一个单列 value 的整数流,用于跑求和示例:

import pathway as pw

# 用人工数据流替代 './sum_input_data/'
t = pw.demo.range_stream(
    nb_rows=30,              # 生成行数,默认 30
    offset=0,                # 加在 value 上的偏移量,默认 0
    input_rate=1.0,          # 每秒生成行数,默认 1.0
    autocommit_duration_ms=1000,  # 提交间隔,默认 1000ms
)
t = t.reduce(pw.reducers.sum(t.value))
pw.io.csv.write(t, "output.csv")
pw.run()

value 列取值范围为 offsetnb_rows + offset(如 nb_rows=50, offset=10 时为 10 到 60)。该函数内部同样走 generate_custom_stream,按 input_rate 匀速产行、每 autocommit_duration_ms 毫秒提交一批——可以完整复现"文件不断新增、结果不断修订"的流式行为。其他合成流(多列、浮点、随机等)见 人工数据流 文档与 demo 模块源码

小结:先批后流的工作流

Pathway 把批到流的切换做到了极简,由此形成一套推荐的工作流:

  1. 先写静态代码:用 mode="static" 的连接器(或静态数据)开发聚合、连接等逻辑,在静态数据上测试、验证正确性;
  2. 切换连接器:把数据源换成支持流式的连接器并把 mode 置为 "streaming"(CSV 连接器本就双模支持),其余管道代码原封不动;
  3. 让框架接管时序:迟到数据、乱序数据、结果修订全部由 Pathway 处理,你需要适配的唯一事实是"输出本身是一条流"。

适用前提与限制:该方案依赖所用连接器同时支持 static 与 streaming 两种模式(以 支持的数据源列表 为准);mode 的默认值是 "streaming",写批处理脚本时务必显式声明 "static"。如果你的流水线原本来自 Pandas,可先看 从 Pandas 迁移 再按本文完成批流切换。

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