首页
/ Pathway 首个实时应用实战:从 CSV 求和到 Kafka 阈值告警的流式 ETL 流水线

Pathway 首个实时应用实战:从 CSV 求和到 Kafka 阈值告警的流式 ETL 流水线

2026-09-06 21:51:06作者:姚月梅Lane

本文以 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),每一行都带有 timediff 字段,这一点在下一节会详细展开;
  • pw.run() 启动计算后,引擎会持续轮询输入端,任何输入变化都会触发整条计算图的增量更新,直到进程被终止——这是 Pathway 的正常行为,而非挂死

如果你希望先用静态、有限的数据测试流水线,可以把连接器改为 mode="static"(一次性读入当前已存在的数据),或者使用官方 demo 模块 构造人工数据流,具体见 流式与静态模式文档 与 artificial streams 文档

示例二:Kafka 实时数据 × CSV 阈值表的 join + filter 告警

这是官方入门文档给出的完整生产级案例,也是理解 Pathway 核心能力的最佳样本:

数据源(两个):

  1. Kafka topic 中的实时测量数据(live measurements);
  2. 本地 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_millistopicpartitionoffset 与 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_diff1-1,即插入/删除标记),均为 UTF-8 字符串。这意味着即使输出到 Kafka 这样的"日志型"系统,下游依然能精确还原每一行的生命周期,与下文 time/diff 列的语义完全一致。

此外 topic_name 不仅可以是字符串常量,还可以是一个字符串列的引用,让每行消息路由到不同 topic;sort_by 参数可以在每个 minibatch 内对输出按指定列升序排序,保证日志的可复现顺序。

pw.run() 与"永远运行"的计算

调用 pw.run() 后,计算图正式启动。之后任何一次输入变化——无论是 Kafka 收到新消息,还是 ./threshold-data/ 下出现新文件、旧文件被修改——都会自动触发整条流水线的增量更新:joined_tablealerts_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}

多出的两行含义清晰:

  1. B 的旧告警被撤回(diff=-1)——因为 B 的新阈值变成了 11,10 不再越限;
  2. 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、逻辑时间等底层模型。

仓库中还有可直接参照的完整项目:

小结

本文沿官方 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 流水线中。

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