Pathway 首个实时应用实战:从 CSV 求和到 Kafka 阈值告警的流式 ETL 流水线
本文以 Pathway Live Data Framework(简称 Pathway)官方入门文档为骨架,带你从零搭建第一个实时 ETL 应用:先跑通"读 CSV、过滤正值、求和、写 JSON Lines"的最小流水线,再实现一个生产级的"Kafka 实时测量数据 × 本地 CSV 阈值表"join + filter 告警系统。读完本文,你能掌握 pw.Schema、输入/输出连接器、pw.run() 计算图的完整用法,并深入理解 Pathway 输出中 time/diff 列背后的"插入-删除日志"(update log)模型。
安装与环境准备
在 Python 3.10+ 环境中,用一条 pip 命令即可安装整个框架(包含 Rust 引擎及运行流水线所需的全部基础依赖):
pip install pathway
适用前提与限制(来自仓库 安装文档 与 pyproject.toml):
- pyproject.toml 中声明
requires-python = ">=3.10",与文档要求的 Python 3.10+ 一致; - Pathway 目前支持 MacOS 与 Linux,暂不支持 Windows(Windows 用户可考虑 WSL、Docker 或虚拟机);
- 标准安装不会引入 LLM 相关库;若后续要跑 Live Data AI 流水线,可另行安装可选依赖组,如
pip install "pathway[xpack-llm]"。
安装完成后,一个 Python 文件就是一个完整的实时应用,无需部署独立的作业系统。
示例一:最小可运行流水线——对 CSV 正值求和
官方入门文档的第一个例子:对 CSV 文件中的正值求和,并将结果写入一个 JSON Lines 文件。这是理解 Pathway 编程模型的最好起点,因为整条流水线只有三个要素——输入连接器 → 表操作 → 输出连接器:
import pathway as pw
from pathlib import Path
class SumsSchema(pw.Schema):
value: float
# 1. 输入:从目录读取所有 CSV,流式模式会自动感知文件的新增/修改
input_table = pw.io.csv(Path("data/"), schema=SumsSchema, mode="streaming")
# 2. 表操作:过滤正值并按列求和(sum 是聚合操作,输出单行结果)
result_table = input_table.filter(pw.this.value > 0).sum(pw.this.value)
# 3. 输出:把求和结果的更新流写入 JSON Lines 文件
pw.io.jsonlines.write(result_table, Path("output.jsonl"))
# 4. 启动计算
pw.run()
对应到源码可以验证几个关键行为:
pw.io.csv实际入口是 csv 连接器,其read函数默认mode="streaming"、autocommit_duration_ms=1500,并且内部委托给更通用的pw.io.fs.read(..., format="csv"),因此文件新增、删除、修改都会被跟踪并反映到表状态中(python/pathway/io/csv/init.py);- 输出侧
pw.io.jsonlines.write会把表的更新流(而不只是当前快照)序列化到文件(python/pathway/io/jsonlines/init.py),每一行都带有time与diff字段,这一点在下一节会详细展开; pw.run()启动计算后,引擎会持续轮询输入端,任何输入变化都会触发整条计算图的增量更新,直到进程被终止——这是 Pathway 的正常行为,而非挂死。
如果你希望先用静态、有限的数据测试流水线,可以把连接器改为 mode="static"(一次性读入当前已存在的数据),或者使用官方 demo 模块 构造人工数据流,具体见 流式与静态模式文档 与 artificial streams 文档。
示例二:Kafka 实时数据 × CSV 阈值表的 join + filter 告警
这是官方入门文档给出的完整生产级案例,也是理解 Pathway 核心能力的最佳样本:
数据源(两个):
- Kafka topic 中的实时测量数据(live measurements);
- 本地 CSV 文件中的阈值配置(thresholds,会随时间被修改)。
目标:将两个数据源按 name 关联,找出当前值超过阈值的实时测量,并把告警写回 Kafka 的另一个 topic。
完整源码
import pathway as pw
# 用 pw.Schema 声明表结构。
# 两张输入表:(1) measurements 是实时流;(2) threshold 是可被修改的 CSV。
# 两者都有两列:name (str) 和一个 float 列。
class MeasurementSchema(pw.Schema):
name: str
value: float
class ThresholdSchema(pw.Schema):
name: str
threshold: float
# Kafka 连接配置(librdkafka 格式)
rdkafka_settings = {
"bootstrap.servers": "server-address:9092",
"security.protocol": "sasl_ssl",
"sasl.mechanism": "SCRAM-SHA-256",
"group.id": "$GROUP_NAME",
"session.timeout.ms": "6000",
"sasl.username": "username",
"sasl.password": "********",
}
# 通过 Kafka 连接器读取实时测量数据
measurements_table = pw.io.kafka.read(
rdkafka_settings,
topic="topic",
schema=MeasurementSchema,
format="json",
autocommit_duration_ms=1000
)
# 通过 CSV 连接器读取阈值文件(目录会被持续监听)
thresholds_table = pw.io.csv(
'./threshold-data/',
schema=ThresholdSchema,
)
# 按 name 列做 join
joined_table = (
# 左表为 measurements_table(记作 pw.left)
measurements_table
.join(
# 右表为 thresholds_table(记作 pw.right)
thresholds_table,
# 两表都按 name 列关联
pw.left.name==pw.right.name,
)
# 用 select 挑选 join 后的输出列
.select(
# 保留 measurements 的全部列
*pw.left,
# 保留 thresholds 表的 threshold 列
pw.right.threshold
)
)
# 过滤出严格超过阈值的记录
alerts_table = (
joined_table
# 仅保留 value 严格大于 threshold 的行
.filter(pw.this.value > pw.this.threshold)
# 输出只保留 name 与 value 两列
.select(pw.this.name, pw.this.value)
)
# 把告警写回同一 Kafka 实例的另一个 topic
pw.io.kafka.write(
alerts_table, rdkafka_settings, topic_name="alerts_topic", format="json"
)
# 启动 Pathway 计算
pw.run()
小注:官方文档原文在 filter 一步写的是
joined_values,结合上下文应为joined_table的笔误;上面的代码已修正为可直接运行的写法。
数据流向
整条流水线的拓扑是:
Kafka topic ──(kafka.read, json)──▶ measurements_table ──┐
├─ join(name) ─ filter(value>threshold) ─▶ alerts_table ─(kafka.write, json)─▶ alerts_topic
./threshold-data/*.csv ──(csv, 目录监听)─▶ thresholds_table ┘
源码级参数解析
结合 Kafka 连接器源码,pw.io.kafka.read 的完整签名与本文用到的参数如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
rdkafka_settings |
必填 | librdkafka 格式的连接配置字典(bootstrap.servers、SASL 认证等) |
topic |
None |
要读取的 topic 名称,支持单个或多个 |
schema |
None |
结果表结构;format="json" 时必填 |
mode |
"streaming" |
streaming 会持续等待新消息;static 只读入执行时刻已存在的数据 |
format |
"raw" |
支持 "plaintext" / "raw" / "json";JSON 模式下按 schema 解析每个消息的字段 |
autocommit_duration_ms |
1500 |
两次 commit 之间的最大间隔(毫秒):每隔该时长,连接器收到的更新会批量提交进计算图。示例中设为 1000 |
json_field_paths |
None |
JSON 模式下把字段名映射到 JSON Pointer 路径(RFC 6901),可提取嵌套字段 |
with_metadata |
False |
为 True 时附加 _metadata 列,含 timestamp_millis、topic、partition、offset 与 headers |
start_from_timestamp_ms |
None |
从指定的历史时间点(毫秒)开始消费 |
parallel_readers |
None |
并行 reader 副本数,未指定时取 min{pathway 线程数, 分区数} |
max_backlog_size |
None |
限制在途处理中的条目数,达到上限时暂停读取,适合源头初始爆发式写入的场景 |
CSV 侧的 pw.io.csv 连接器 同样默认 mode="streaming"、autocommit_duration_ms=1500:传入目录路径 ./threshold-data/ 后,引擎会按文件修改时间排序处理目录内文件,并持续监听文件的新增、删除与修改——删除文件会从表中移除对应行,这正是示例二"阈值文件更新后告警自动重算"的底层机制。
输出侧,pw.io.kafka.write 支持 "json"、"dsv"、"plaintext"、"raw" 四种序列化格式。有一个值得注意的实现细节:写出的每条 Kafka 消息除了 key/value 外,还会附带两个 headers——pathway_time(该条目的逻辑时间)与 pathway_diff(1 或 -1,即插入/删除标记),均为 UTF-8 字符串。这意味着即使输出到 Kafka 这样的"日志型"系统,下游依然能精确还原每一行的生命周期,与下文 time/diff 列的语义完全一致。
此外 topic_name 不仅可以是字符串常量,还可以是一个字符串列的引用,让每行消息路由到不同 topic;sort_by 参数可以在每个 minibatch 内对输出按指定列升序排序,保证日志的可复现顺序。
pw.run() 与"永远运行"的计算
调用 pw.run() 后,计算图正式启动。之后任何一次输入变化——无论是 Kafka 收到新消息,还是 ./threshold-data/ 下出现新文件、旧文件被修改——都会自动触发整条流水线的增量更新:joined_table 与 alerts_table 随之刷新,变化经 Kafka 输出连接器转发到 alerts_topic。
引擎会不断轮询新的更新,直到进程被终止才停止。官方文档特别强调:"This is the normal behavior of the framework"——实时应用的进程常驻轮询就是设计意图,不需要额外的调度器。
深入理解输出:time 与 diff 列(更新日志模型)
Pathway 的输出不是"当前状态的快照",而是一份插入/删除操作日志(log of insertion and suppression),每个输出的行都额外携带两个字段:
time:该行更新发生的时刻(逻辑时间,实践中对应常规时间戳的推进);diff:标识该行是插入(diff = 1)还是删除(diff = -1)。一次"更新"由两行表示:一行删除旧值、一行插入新值,两行共享同一个time,以保证操作的原子性。
场景推演
假设 Kafka topic 收到如下输入:
{"name": "A", "value":8}
{"name": "B", "value":10}
而阈值文件内容为:
name, threshold
"A", 9
"B", 9
只有 B(10 > 9)超过阈值,输出为:
{"name": "B", "value":10, "time":1, "diff":1}
这里假设首批值在 time=1 时计算完成,diff=1 表示插入。
现在阈值文件被修改为:
name, threshold
"A", 7
"B", 11
CSV 连接器会自动检测 ./threshold-data/ 下的变更并更新 thresholds_table,进而触发 join 与 filter 重新执行,输出追加为:
{"name": "B", "value":10, "time":1, "diff":1}
{"name": "B", "value":10, "time":2, "diff":-1}
{"name": "A", "value":8, "time":2, "diff":1}
多出的两行含义清晰:
- B 的旧告警被撤回(
diff=-1)——因为 B 的新阈值变成了 11,10 不再越限; - A 产生新告警(
diff=1)——因为 A 的阈值从 9 降到 7,8 现在越限了。
注意旧的行依然保留在输出中:这正是更新日志模型的价值——它提供了关于数据"发生过什么"的完整信息。同时官方文档提醒:某些面向外部存储的输出连接器会把这些 +1/-1 行按 time 配对、合并表示成一次 update 操作(例如写入 SQL 数据库时表现为 UPDATE 而非 DELETE+INSERT),无需你手工处理。
相关文档与更多示例
- 连接器总览:live-data-framework-connectors 汇总了所有输入连接器;Kafka 连接器专页 覆盖更多认证与配置场景。
- Schema 定义:schema 文档 讲解
pw.Schema的完整用法,同一个 Schema 类可以被多张表复用。 - 表操作:table operations 指南 展示 join、时间窗口、filter、groupby 等可用操作。
- 流式 vs 静态模式:streaming-and-static-modes 说明如何用
mode="static"与pw.demo人工数据流对流水线做静态测试。 - 核心概念:concepts 专文 系统解释 minibatch、逻辑时间等底层模型。
仓库中还有可直接参照的完整项目:
- examples/projects/kafka-ETL:Kafka 实时 ETL 完整工程;
- examples/projects/realtime-log-monitoring:事件驱动 + 告警的实时日志监控流水线;
- examples/projects/kafka-linear-regression:Kafka 实时分析(线性回归)。
小结
本文沿官方 first-realtime-app 文档 的脉络走完了 Pathway 的首个实时应用:一条 pip 命令安装(Python 3.10+,MacOS/Linux);用三行核心代码搭起"CSV → 过滤求和 → JSON Lines"的最小流水线;用约 70 行代码实现 Kafka 实时测量与 CSV 阈值表的 join + filter 告警,并将结果写回 Kafka。关键在于理解两点:其一,pw.run() 启动后整条计算图常驻轮询、随输入增量更新,这是实时应用的常态而非异常;其二,所有输出都是带 time/diff 的插入-删除日志,它既保证了更新原子性,也让你能完整追溯每一条数据的生命周期。掌握这两点后,你就可以把同样的模式推广到日志监控、实时分析与 Live Data AI 流水线中。
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 StartedRust0624
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