Pathway 持久化实战:从流式词频统计到中断计算的断点恢复
本文基于 Pathway 官方教程,完整走一遍"无持久化 → 有持久化"的流式计算恢复流程:先构建一个持续轮询输入目录的流式词频统计(word count)程序,验证其在进程被终止后重启会从头重算;再仅用两行持久化配置加一个数据源命名,让同一个程序在重启后从上次停下的位置继续计算。读完本文,你将掌握 pw.persistence.Backend / pw.persistence.Config 的完整配置项、pw.run 中 persistence_config 的接入方式,以及理解 Pathway "at-least-once" 重放语义的边界。
示例任务:流式目录上的词频统计
整个实验围绕一个基础词频统计任务展开:
- Pathway 程序持续扫描文件系统中的一个目录(
inputs/),目录内容是若干 CSV 文件; - 每个 CSV 文件的格式固定为:一列名为
word,每个文件包含恰好一个单词; - 输出为 JsonLines 文件(
result.jsonlines),其中流式记录两个字段的变更:单词word及其计数count。
程序以独立的流式进程方式运行(即一直等待输入),另一个进程充当 streamer,每隔一段时间向输入目录投放一个新文件。
环境准备:输入目录清理与数据投放器
为了对比"无持久化"和"有持久化"两个变体,先准备两个辅助方法。
1. 清理输入目录。 输入统一存放在 inputs/ 目录,每次实验前必须保证它是空的:
import multiprocessing
import os
import shutil
import subprocess
import time
def clean_input_directory():
if os.path.exists("inputs/"):
shutil.rmtree("inputs/")
os.mkdir("inputs/")
clean_input_directory()
2. 投放器(streamer)。 它的目标是周期性地在输入目录中生成新文件。接口上做了三点设计:
- 通过
interval_sec参数控制新文件的投放间隔; - 通过
how_many参数指定投放的输入块(文件)总数; - 词本身采用轮询(round-robin)方式从
words列表中生成,保证可复现。
def generate_inputs(interval_sec, how_many, words):
for file_id in range(how_many):
time.sleep(interval_sec)
with open(f"inputs/{time.time()}", "w") as f:
f.write(f"word\n{words[file_id % len(words)]}")
注意 time.time() 在多数平台上具有微秒级精度,因此用它做文件名可以方便地保证文件名的唯一性——这也是 CSV 连接器按文件修改时间升序处理目录内文件(见 csv.read 的文档)时能正确排序的关键。
第一步:不启用持久化的 Pathway 程序
现在编写一个不含持久化的词频统计程序,分三个步骤:
- 定义输入连接器:数据来自 CSV 目录,使用
pw.io.csv.read,源指向inputs/,传入 schema,并把自动提交时间设为 10 毫秒(autocommit_duration_ms=10),确保更新频繁地推入引擎。该参数的默认值是 1500 毫秒(见 read 签名),调小它能加快变更的可见性; - 定义计算逻辑:词频统计本质上是对单词流的 group-by 操作,用
groupby加reduce,reducer 统计每个词的出现次数并挂到该词上; - 定义输出连接器:结果以 JsonLines 格式写入文件系统,使用
pw.io.jsonlines.write。
另外两个工程细节:为减少干扰,通过向 pw.run 传递关键字参数关闭 top 风格的监控视图(monitoring_level=pw.MonitoringLevel.NONE);为模拟"应用作为独立进程运行"这一真实部署形态,把计算代码存到单独的 wordcount.py,再用 subprocess.Popen 启动它。
程序代码如下:
import pathway as pw
class InputSchema(pw.Schema):
word: str
if __name__ == "__main__":
words = pw.io.csv.read(
"inputs/",
schema=InputSchema,
autocommit_duration_ms=10
)
word_counts = words.groupby(words.word).reduce(words.word, count=pw.reducers.count())
pw.io.jsonlines.write(word_counts, "result.jsonlines")
pw.run(monitoring_level=pw.MonitoringLevel.NONE)
启动逻辑非常简单:
def run_pathway_wordcount_program():
pd = subprocess.Popen(["python", "wordcount.py"])
return pd
实验:中断后重启,计算从头开始
现在联合测试投放器与程序:以 50 毫秒的间隔生成 200 个输入文件,文件内容在 "hello" 与 "world" 之间交替。10 秒后共 200 个文件,其中各 100 个包含 "hello" 和 "world"。
启动流式投放后立即启动 Pathway 程序,等待 5 秒(此时投放器尚未产出全部文件)就用 subprocess.Popen 的 kill 方法杀掉它,再等 5 秒让投放器完成:
# Start streaming inputs
pd_gen = multiprocessing.Process(
target=generate_inputs,
args=(0.05, 200, ["hello", "world"])
)
pd_gen.start()
# Run Pathway program
pd_comp = run_pathway_wordcount_program()
time.sleep(5)
pd_comp.kill()
pd_gen.join()
查看结果文件尾部:
!tail -5 result.jsonlines
{"word":"hello","count":49,"diff":1,"time":1699279664772}
{"word":"world","count":48,"diff":-1,"time":1699279664822}
{"word":"world","count":49,"diff":1,"time":1699279664822}
{"word":"hello","count":49,"diff":-1,"time":1699279664874}
{"word":"hello","count":50,"diff":1,"time":1699279664874}
结果显然是不完整的:每个词只计到 50,而实际应各有 100 个。原因就是程序在正常执行期间被强制终止了。
接着直接重新运行程序(此时目录中所有文件已就位):
pd_comp = run_pathway_wordcount_program()
time.sleep(5)
pd_comp.kill()
!tail -5 result.jsonlines
{"word":"hello","count":98,"diff":1,"time":1699279678510}
{"word":"world","count":97,"diff":-1,"time":1699279678512}
{"word":"world","count":100,"diff":1,"time":1699279678512}
{"word":"hello","count":98,"diff":-1,"time":1699279678512}
{"word":"hello","count":100,"diff":1,"time":1699279678512}
这次最终计数达到了预期的 100/100。但再看输出文件的开头:
!head -5 result.jsonlines
{"word":"world","count":3,"diff":1,"time":1699279678436}
{"word":"hello","count":4,"diff":1,"time":1699279678436}
{"word":"world","count":3,"diff":-1,"time":1699279678438}
{"word":"world","count":4,"diff":1,"time":1699279678438}
{"word":"hello","count":4,"diff":-1,"time":1699279678438}
可以看出,程序从计数 1 开始从头计算——无持久化的程序在重启后必然把全部输入重读、重算、重写一遍。那么如何避免这件事?
引入持久化:原理与后端配置
持久化(Persistence)是让程序记住上次执行时"计算读到哪、输出到哪"的机制。其核心思想是:框架周期性地把计算状态 dump 到指定的数据存储后端;重启时,框架先在后端寻找快照,找到则把快照加载进引擎,从而不需要对已保存的数据重新做读取和处理。
持久化后端
持久化机制保存的计算快照由两部分组成:一部分是与输入规模大致成正比的原始数据,另一部分是体量较小的元数据。这两者都必须落在持久性存储上。Pathway 目前提供的存储选择包括:
- 文件系统,即本地磁盘上的一个目录;
- S3 桶(可指定根目录);
- 从 Backend 源码 看,还实现了 S3 与 Azure Blob Storage 之外的
azure类方法(Backend.azure(root_path, account, password, container)),教程示例只演示了其中一种。
本实验使用本地文件系统的持久化存储,假设放在名为 PStorage 的目录中。由于实验可能需要反复运行,先准备一个清理辅助方法:
def clean_persistent_storage():
if os.path.exists("./PStorage"):
shutil.rmtree("./PStorage")
clean_persistent_storage()
然后是持久化配置本身,共两行:
backend = pw.persistence.Backend.filesystem("./PStorage")
persistence_config = pw.persistence.Config(backend)
第一行用 pw.persistence.Backend.filesystem 创建文件系统后端(对应 Backend.filesystem 的实现,它只需一个存放持久化数据的根目录路径);第二行用 pw.persistence.Config 构造器创建持久化配置,唯一的位置参数就是刚创建的后端。
pw.persistence.Config 是一个 frozen dataclass,除 backend 外还暴露了若干可选参数,可从 源码定义 得到完整取值:
| 参数 | 默认值 | 含义 |
|---|---|---|
backend |
必填 | 持久化后端配置 |
snapshot_interval_ms |
0 |
期望的快照更新间隔(毫秒)。值越大,快照落后越多,但计算资源消耗越低 |
snapshot_access |
api.SnapshotAccess.FULL |
快照访问级别 |
persistence_mode |
api.PersistenceMode.PERSISTING |
持久化模式:PERSISTING(默认,持久化全部数据)、UDF_CACHING(仅缓存 UDF 调用结果)、OPERATOR_PERSISTING(仅持久化内部算子状态,最省开销) |
continue_after_replay |
True |
重放后是否继续 |
worker_scaling_enabled |
False |
动态 worker 扩缩容,启用时要求程序通过 pathway spawn 启动 |
workload_tracking_window_ms |
120000 |
扩缩容时评估管道负载的时间窗口 |
Python 侧的字段定义与 Rust 引擎侧的构造签名一一对应,可参见 src/python_api.rs 中的 PersistenceConfig。这些配置最终通过 pw.run 的 persistence_config 参数传入引擎——run 函数签名 中明确列出了 persistence_config: PersistenceConfig | None = None,默认不启用持久化,因此非持久程序与持久程序的计算图代码可以完全相同,差别只在 pw.run 的入参。
仓库中的 test_persistence.py 正是以同一个"CSV 词频统计 + snapshot_interval_ms=1000"场景对持久化做了自动化回归验证,可以作为本教程示例的测试侧印证。
唯一命名(Unique Names)
第二件(可选的)事情是为数据源分配唯一名字。唯一名是引擎在不同运行之间匹配数据源所必需的。
原则上引擎可以自动完成这一分配——按数据源出现和构造的顺序依次命名。但如果未来你会修改 Pathway 程序或增删数据源,这种隐式命名就不够稳妥,官方并不推荐。为完整起见,本教程演示手动命名:与无持久化版本的输入相比,唯一区别就是给 pw.io.csv.read 多传一个 name 参数。例如把数据源命名为 words_data_source:
pw.io.csv.read(
...,
name="words_data_source"
)
name 参数的语义在 csv.read 的文档字符串 中有明确说明:若提供该名字,它会被用于日志与监控看板;当持久化启用时,它同时作为存储连接器进度的快照的名字——这正是重启后引擎能把"上次的读取进度"与"本次的连接器"对应起来的关键。
另外要注意名字必须全局唯一:从 test_io.py 中的测试 可以看到,若输入连接器和输出连接器使用了相同的 name,pw.run 会直接抛出 ValueError: Unique name 'one' used more than once。
重新审视应用:两行改动换来断点续算
把上述改动应用到最初的程序上。接口保持不变——仍然是可以被打断的独立进程;持久化版本的 wordcount.py 与原版只有两处差异(代码中已用注释标出):
import pathway as pw
class InputSchema(pw.Schema):
word: str
if __name__ == "__main__":
words = pw.io.csv.read(
"inputs/",
schema=InputSchema,
autocommit_duration_ms=10,
name="words_input_source", # Changed: now name is assigned here
)
word_counts = words.groupby(words.word).reduce(words.word, count=pw.reducers.count())
pw.io.jsonlines.write(word_counts, "result.jsonlines")
backend = pw.persistence.Backend.filesystem("./PStorage")
persistence_config = pw.persistence.Config(backend)
pw.run(
monitoring_level=pw.MonitoringLevel.NONE,
persistence_config=persistence_config, # Changed: now persistence_config is passed here
)
实验:中断,再重启
和上次一样:以 50 毫秒间隔生成 200 个交替包含 "hello" / "world" 的文件,并在投放完成前终止程序:
# Clean the old files: remove old results and inputs
!rm -rf result.jsonlines
clean_input_directory()
# Start streaming inputs
pd_gen = multiprocessing.Process(
target=generate_inputs,
args=(0.05, 200, ["hello", "world"])
)
pd_gen.start()
# Run Pathway program
pd_comp = run_pathway_wordcount_program()
time.sleep(5)
pd_comp.kill()
检查输出,此时程序最多消费了约一半的输入:
!tail -5 result.jsonlines
{"word":"hello","count":49,"diff":1,"time":1699279708352}
{"word":"world","count":48,"diff":-1,"time":1699279708402}
{"word":"world","count":49,"diff":1,"time":1699279708402}
{"word":"hello","count":49,"diff":-1,"time":1699279708452}
{"word":"hello","count":50,"diff":1,"time":1699279708452}
现在在全量输入就位的情况下再次运行程序(流式模式下它不会自行结束,因此运行 5 秒后终止):
pd_comp = run_pathway_wordcount_program()
time.sleep(5)
pd_comp.kill()
检查结果时,尾部用于验证结果正确性(每个词的最终计数应为 100),头部用于观察程序从何处开始产生输出:
!head -5 result.jsonlines
!echo "==="
!tail -5 result.jsonlines
{"word":"world","count":49,"diff":-1,"time":1699279716584}
{"word":"world","count":51,"diff":1,"time":1699279716584}
{"word":"hello","count":50,"diff":-1,"time":1699279716584}
{"word":"hello","count":51,"diff":1,"time":1699279716584}
{"word":"world","count":51,"diff":-1,"time":1699279716586}
===
{"word":"hello","count":99,"diff":1,"time":1699279716634}
{"word":"world","count":98,"diff":-1,"time":1699279716636}
{"word":"world","count":100,"diff":1,"time":1699279716636}
{"word":"hello","count":99,"diff":-1,"time":1699279716636}
{"word":"hello","count":100,"diff":1,"time":1699279716636}
两点观察:
- 结果正确——两个词最终都计到了 100;
- 新产生的输出从计数 50 附近(具体位置随运行而异)开始,说明程序没有重算、重写此前已计算并输出的数据,而是从上次停止处继续。
语义边界:at-least-once 重放
需要注意,上面输出的"前几行"可能与上一次运行的"后几行"发生重叠。这里体现的是 at-least-once 语义:在初始计算被中断时,有一个事务迷你批(transaction minibatch)尚未提交,重启后它会被重复投递一次。对词频统计这类可交换、可重入的聚合,重复投递不影响最终结果的正确性;但如果你在下游写自定义副作用逻辑,应意识到恢复点并非严格"恰好一次"。
另外从 Config 的源码 可以看到,配置对象在 pw.run 前后会执行 on_before_run / on_after_run,把文件系统后端的根路径写入(并随后清理)环境变量 PATHWAY_PERSISTENT_STORAGE——仓库内大量测试(如 test_io.py)正是通过这个环境变量来搭建持久化测试环境的,这解释了持久化状态目录与运行进程之间的耦合方式。
小结
本教程在一个刻意保持简单的词频统计例子上演示了持久化的完整闭环:
- 无持久化:程序重启后从计数 1 开始,全部输入重读重算;
- 有持久化:只需
Backend.filesystem+Config(backend)两行配置,并通过name参数给连接器固定快照名,重启即从上次停止处继续; - 语义上是 at-least-once:未提交的事务迷你批在恢复时可能被重复投递,最终聚合结果仍正确。
持久化是一个更广义的机制,除断点恢复外还能处理其他任务——例如在满足一定条件时处理数据源变更,参见同目录下的 数据源变更后重启教程;持久化的核心概念(持久化存储、快照与写前日志)则有更系统的阐述,见 持久化概念文档。
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