
Pathway 首个实时应用实战从 CSV 求和到 Kafka 阈值告警的流式 ETL 流水线【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文以 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.tomlpyproject.toml 中声明requires-python 3.10与文档要求的 Python 3.10 一致Pathway 目前支持MacOS 与 Linux暂不支持 WindowsWindows 用户可考虑 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/), schemaSumsSchema, modestreaming) # 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函数默认modestreaming、autocommit_duration_ms1500并且内部委托给更通用的pw.io.fs.read(..., formatcsv)因此文件新增、删除、修改都会被跟踪并反映到表状态中python/pathway/io/csv/init.py输出侧pw.io.jsonlines.write会把表的更新流而不只是当前快照序列化到文件python/pathway/io/jsonlines/init.py每一行都带有time与diff字段这一点在下一节会详细展开pw.run()启动计算后引擎会持续轮询输入端任何输入变化都会触发整条计算图的增量更新直到进程被终止——这是 Pathway 的正常行为而非挂死。如果你希望先用静态、有限的数据测试流水线可以把连接器改为modestatic一次性读入当前已存在的数据或者使用官方 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, topictopic, schemaMeasurementSchema, formatjson, autocommit_duration_ms1000 ) # 通过 CSV 连接器读取阈值文件目录会被持续监听 thresholds_table pw.io.csv( ./threshold-data/, schemaThresholdSchema, ) # 按 name 列做 join joined_table ( # 左表为 measurements_table记作 pw.left measurements_table .join( # 右表为 thresholds_table记作 pw.right thresholds_table, # 两表都按 name 列关联 pw.left.namepw.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_namealerts_topic, formatjson ) # 启动 Pathway 计算 pw.run()小注官方文档原文在 filter 一步写的是joined_values结合上下文应为joined_table的笔误上面的代码已修正为可直接运行的写法。数据流向整条流水线的拓扑是Kafka topic ──(kafka.read, json)──▶ measurements_table ──┐ ├─ join(name) ─ filter(valuethreshold) ─▶ alerts_table ─(kafka.write, json)─▶ alerts_topic ./threshold-data/*.csv ──(csv, 目录监听)─▶ thresholds_table ┘源码级参数解析结合 Kafka 连接器源码pw.io.kafka.read的完整签名与本文用到的参数如下参数默认值说明rdkafka_settings必填librdkafka 格式的连接配置字典bootstrap.servers、SASL 认证等topicNone要读取的 topic 名称支持单个或多个schemaNone结果表结构formatjson时必填modestreamingstreaming会持续等待新消息static只读入执行时刻已存在的数据formatraw支持plaintext/raw/jsonJSON 模式下按 schema 解析每个消息的字段autocommit_duration_ms1500两次 commit 之间的最大间隔毫秒每隔该时长连接器收到的更新会批量提交进计算图。示例中设为1000json_field_pathsNoneJSON 模式下把字段名映射到 JSON Pointer 路径RFC 6901可提取嵌套字段with_metadataFalse为True时附加_metadata列含timestamp_millis、topic、partition、offset与 headersstart_from_timestamp_msNone从指定的历史时间点毫秒开始消费parallel_readersNone并行 reader 副本数未指定时取min{pathway 线程数, 分区数}max_backlog_sizeNone限制在途处理中的条目数达到上限时暂停读取适合源头初始爆发式写入的场景CSV 侧的 pw.io.csv 连接器 同样默认modestreaming、autocommit_duration_ms1500传入目录路径./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不仅可以是字符串常量还可以是一个字符串列的引用让每行消息路由到不同 topicsort_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只有 B10 9超过阈值输出为{name: B, value:10, time:1, diff:1}这里假设首批值在 time1 时计算完成diff1表示插入。现在阈值文件被修改为name, threshold A, 7 B, 11CSV 连接器会自动检测./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 的新阈值变成了 1110 不再越限A 产生新告警diff1——因为 A 的阈值从 9 降到 78 现在越限了。注意旧的行依然保留在输出中这正是更新日志模型的价值——它提供了关于数据发生过什么的完整信息。同时官方文档提醒某些面向外部存储的输出连接器会把这些1/-1行按time配对、合并表示成一次 update 操作例如写入 SQL 数据库时表现为 UPDATE 而非 DELETEINSERT无需你手工处理。相关文档与更多示例连接器总览live-data-framework-connectors 汇总了所有输入连接器Kafka 连接器专页 覆盖更多认证与配置场景。Schema 定义schema 文档 讲解pw.Schema的完整用法同一个 Schema 类可以被多张表复用。表操作table operations 指南 展示 join、时间窗口、filter、groupby 等可用操作。流式 vs 静态模式streaming-and-static-modes 说明如何用modestatic与pw.demo人工数据流对流水线做静态测试。核心概念concepts 专文 系统解释 minibatch、逻辑时间等底层模型。仓库中还有可直接参照的完整项目examples/projects/kafka-ETLKafka 实时 ETL 完整工程examples/projects/realtime-log-monitoring事件驱动 告警的实时日志监控流水线examples/projects/kafka-linear-regressionKafka 实时分析线性回归。小结本文沿官方 first-realtime-app 文档 的脉络走完了 Pathway 的首个实时应用一条 pip 命令安装Python 3.10MacOS/Linux用三行核心代码搭起CSV → 过滤求和 → JSON Lines的最小流水线用约 70 行代码实现 Kafka 实时测量与 CSV 阈值表的 join filter 告警并将结果写回 Kafka。关键在于理解两点其一pw.run()启动后整条计算图常驻轮询、随输入增量更新这是实时应用的常态而非异常其二所有输出都是带time/diff的插入-删除日志它既保证了更新原子性也让你能完整追溯每一条数据的生命周期。掌握这两点后你就可以把同样的模式推广到日志监控、实时分析与 Live Data AI 流水线中。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考