首页
/ Pathway 持久化机制详解:持久存储后端、快照恢复与至多/至少一次的交付语义

Pathway 持久化机制详解:持久存储后端、快照恢复与至多/至少一次的交付语义

2026-09-04 14:01:26作者:盛欣凯Ernestine

本文围绕 Pathway Live Data Framework 的持久化(Persistence)机制展开:以官方向导中的 wordcount 示例为骨架,完整讲解如何通过 pw.persistence.Backendpw.persistence.Config 两行配置让流式程序在重启后从上次的断点继续计算;并结合仓库源码深入剖析快照、元数据存储、name 唯一标识以及 at-least-once 交付语义的底层实现,帮助你在生产环境中可靠地恢复被中断的数据管道。

一、为什么流式管道需要持久化

Pathway 帮助开发者以声明式方式构建计算管道(computation pipeline)。在开发和运行过程中,经常需要保存计算的状态,典型动机有两个(见官方文档 55.persistence.md):

  1. 状态保存:程序重启后能从上次中断的位置继续工作,而不是从头再来;
  2. 故障恢复:从一次失败中恢复时,不必把整个数据管道从头重新计算一遍。

Pathway 的持久化机制将“内部状态序列化”与“预写日志(write-ahead logging)”结合起来:引擎周期性地把计算状态转储(dump)到持久存储后端;重启时先查找已持久化的 checkpoint,加载快照数据以及各数据源中下一条未读 offset 的信息,从而跳过已处理过的数据。

二、从一个非持久化的 wordcount 开始

官方文档给出的示例任务:输入是一个持续被轮询(polling)的 CSV 文件目录,每个文件含一行表头和若干行单词;输出是一个 JSON Lines 文件,每行包含 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)

经过以上修改,程序即成为持久化程序:重启后会从上次停止处继续计算,并继续产出输出。pw.runpersistence_config 参数的定义见 python/pathway/internals/run.py,其中 docstring 明确说明该参数用于“在需要持久化时保存状态”(the config for persisting the state in case this persistence is required)。

官方仓库还提供了配套的分步教程 60.persistence_recovery.md:它用一个“streamer”脚本不断向输入目录写入 CSV 文件,中途 kill 掉 Pathway 进程,再对比“无持久化”(重启后从头计数)与“有持久化”(重启后从约 50 的计数继续)两种运行结果,直观验证了断点恢复行为。

三、持久化配置的完整参数:pw.persistence.Config

pw.persistence.Config 实例会累积多项设置,并作为参数传给 pw.run。官方文档概括了其中最核心的三项:元数据存储(metadata storage)快照存储(snapshot storage)快照间隔(snapshot interval)

  • 元数据存储:保存恢复所需的小型元信息——时间推进(times advanced)、当前各数据源的位置、计算图的描述等。其大小与处理的数据量无关;
  • 快照存储:保存数据快照(data snapshot),是更大的结构体,大小取决于数据量;
  • 快照间隔:期望的快照新鲜度(freshness)。它是“占用过多计算资源”与“快照不够新”之间的权衡——决定了新更新距离最后一次落盘状态可以有多近,引擎才允许暂缓存储它们。

从源码看(python/pathway/persistence/init.py),Config 是一个 frozen dataclass,完整字段及默认值如下:

参数 类型 / 默认值 含义
backend pw.persistence.Backend(必填) 持久化后端配置
snapshot_interval_ms int = 0 快照更新之间的期望时长(毫秒);值越大,快照允许落后越多,所需计算资源越少
snapshot_access api.SnapshotAccess.FULL 快照的访问方式(如 RECORD / REPLAY)
persistence_mode api.PersistenceMode.PERSISTING 持久化模式,见下文三种取值
continue_after_replay bool = True replay 结束后是否继续处理
worker_scaling_enabled bool = False 启用 worker 进程动态伸缩;注意该能力要求程序以 pathway spawn 方式启动
workload_tracking_window_ms int = 120000 动态伸缩时评估负载的时间窗口;过载/过闲状态需持续整个窗口才会触发伸缩决策

其中 persistence_mode 在源码 docstring 中给出了三种模式(python/pathway/persistence/init.py):

  • pw.PersistenceMode.PERSISTING:默认值,全部数据均被持久化。指定它(或省略该参数)并把配置传给 pw.run 后,无需任何额外操作即可持久化程序状态;
  • pw.PersistenceMode.UDF_CACHING:只缓存用户自定义函数(UDF)调用。缓存保存“函数入参 → 结果”的映射,再次以相同入参调用时直接返回缓存结果;
  • pw.PersistenceMode.OPERATOR_PERSISTING:最高效的持久化机制,只持久化内部算子的状态,既不保存输入也不对输入做任何重算。

需要提醒的是,仓库中还存在一个旧式构造入口 Config.simple_config(...),源码已明确标注其为 deprecated,仅保留以兼容旧代码,官方建议直接使用 pw.persistence.Config 构造函数(python/pathway/persistence/init.py)。

四、后端选择:pw.persistence.Backend

元数据存储与快照存储都通过 pw.persistence.Backend 配置。官方文档提到它提供 S3 与 Filesystem 两种配置方式;从当前仓库源码看,Backend 实际提供了三个面向用户的类方法(外加一个测试用 mock):

# 本地文件系统后端:path 为持久化数据的根目录
backend = pw.persistence.Backend.filesystem("./state/")

# S3 后端:root_path 为 S3 内的根路径,bucket_settings 与 S3 连接器格式一致
backend = pw.persistence.Backend.s3(root_path, bucket_settings)

# Azure Blob Storage 后端
backend = pw.persistence.Backend.azure(root_path, account, password, container)

定义见 python/pathway/persistence/init.py。在引擎侧(Rust),这些后端被映射为 PersistentStorageConfig 枚举——Filesystem / S3 / Azure / Mock,并分别创建 FilesystemKVStorageS3KVStorageAzureKVStorageMockKVStorage 四类 PersistenceBackend 实现(src/persistence/config.rs),对应的后端实现文件位于 src/persistence/backends/ 目录(file.rss3.rsazure.rsmock.rs)。

一个实现细节值得注意:Config 通过 on_before_run / on_after_run 钩子在运行前把文件系统后端路径写入环境变量 PATHWAY_PERSISTENT_STORAGE、运行结束后删除(python/pathway/persistence/init.py)。pw.run 内部通过 get_persistence_engine_config 上下文管理器统一执行这一“注入配置—执行—清理”流程(python/pathway/persistence/init.py),因此调用者无需关心环境变量生命周期。

五、唯一名称(Unique Names):跨重启识别数据源

要让某个输入源被持久化,框架依赖输入连接器上的 name 参数。这个标识符代表一个“事实上的数据源”,因此预期在多次运行之间保持不变

其动机是:只要数据源包含相同 schema 的数据,它本身可以有任意变化,例如:

  • 数据格式变了:原来用 JSON,新条目改成了 CSV;
  • 路径变了:Pathway 解析的日志现在存放在另一个卷上;
  • 字段被重命名:如 date 改名为 datetime 以表意更精确。

变化可以多种多样,但只要使用相同的 name,引擎就知道这些数据仍然对应同一张表。

name 的分配有两种方式:

  1. 自动生成:按数据源被添加的顺序分配唯一 ID。例如程序先读 CSV 数据集、再读 Kafka 事件流,会自动生成两个唯一名称,第一个指向数据集、第二个指向事件流。这种方式在代码不会改动时没问题;但一旦后续修改代码导致数据源顺序变化,生成的名称就会与旧名称对不上,从而破坏恢复;
  2. 手动指定:在输入连接器中显式传入字符串参数 name,更灵活、推荐用于需要长期演进的生产程序。上文的 wordcount 示例可以改写为:
words = pw.io.csv.read("inputs/", schema=InputSchema, name="words_source")

从源码结构看,name 参数在 CSV 连接器 docstring 中的定义与持久化直接相关:它不仅用于日志和监控看板,而且“如果启用了持久化,它将用作保存连接器进度的快照的名称”(python/pathway/io/csv/init.py)。这正是“同一 name = 同一张表 = 同一份可恢复进度”的实现基础。

六、恢复流程与交付语义:at-least-once 的保证边界

框架在运行中维护内部状态快照及恢复所需的元数据。程序启动时的恢复流程是:

  1. 先查找已持久化的 checkpoint;
  2. 找到后加载快照数据,并读取各数据源中“下一条未读 offset”的信息;
  3. 从断点继续,而不重放已处理过的数据。

由此产生一条硬性前提:数据源本身必须是持久的。这样无论程序因何种原因终止,重启后都能把未读条目重新读入。好消息是多数数据源都满足该要求——S3、Kafka topic、文件系统条目在程序重启后都可以被重新读取。

关于交付语义,官方文档给出了明确边界:

  • 持久化恢复提供 at-least-once(至少一次)交付保证。内部输入会被切分成更小的事务批次(transactional batches);如果程序在运行中被强制中断,对应于“未闭合事务批次”的输出可能会重复出现;
  • 优雅终止(graceful termination)下可以保证 exactly-once(精确一次)语义

这一语义在官方教程的实测中也有体现:中断后重启,输出的开头几行可能与上一次运行结尾相交,出现重复投递,教程将其归因于“初始计算被中断时尚未提交的事务 mini-batch”(60.persistence_recovery.md)。因此若下游要求严格去重,需要在输出端对 at-least-once 的重复做幂等处理;只有在能确保优雅退出的运维条件下,才能依赖 exactly-once 语义。

七、引擎侧实现速览:状态如何被保存与恢复

对想在源码层面理解恢复机制的读者,关键入口位于 Rust 引擎的持久化模块 src/persistence/

Python 侧的测试用例 python/pathway/tests/test_persistence.pypython/pathway/tests/test_persistence_iterate.py 覆盖了持久化与 iterate 模式的组合行为,可作为回归验证的参考。

八、实操清单与注意事项

结合文档与源码,把持久化落到生产时建议核对以下几点:

  1. name 必须稳定。一旦为输入源指定了 name(尤其依赖自动生成时),不要改变数据源的添加顺序,否则恢复会失配——这是文档中明确警告的坑;
  2. snapshot_interval_ms 是资源与新鲜度的权衡。默认值为 0;调大它可降低快照开销,但重启时可能需要重放更多的尾部更新;
  3. 数据源需可重读。S3、Kafka topic、文件系统均满足;如果自研输入源,请确保重启后未读条目可再读取;
  4. 动态伸缩有前提worker_scaling_enabled=True 要求程序通过 pathway spawn 启动,且负载状态需在 workload_tracking_window_ms(默认 120000 ms)窗口内持续存在才会触发伸缩;
  5. 交付语义要有预期。被 kill / 崩溃场景是 at-least-once,优雅退出才是 exactly-once;下游设计需按 at-least-once 做幂等。

官方完整的持久化 API 参考可查阅仓库内的 API 文档目录(docs/2.developers/),本指南所述配置面均以当前仓库 python/pathway/persistence/init.pysrc/persistence/config.rs 的实际实现为准。

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

项目优选

收起
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.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.79 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
988
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384