最近这几周 GitHub Trending 上有个署名很有意思的项目,叫 Pathway。作为一个常年用 Python 写数据处理、又天天被“Python 太慢、实时处理得上 Java/Scala 系”这种话教育的人,我一开始看到标题里“性能吊打 Flink、Spark”这种表述,第一反应是“又是一个标题党”。但真把它拉下来跑完一个 Demo 之后,我的评价变成了:这个项目确实重新定义了我对“Python 能做实时 ETL”这件事的认知。
这篇不是纯安利,也不是劝退文。我会先拆清楚 Pathway 的底层逻辑和性能来源,再看它凭什么敢跟 Flink、Streaming 叫板,然后给你一套可以直接在本地复现的订单流聚合实时 ETL 示例,最后聊聊实际使用中我踩到的坑,以及什么场景下该用它、什么场景下别碰它。
1. Pathway 是什么:用 Python 写实时 ETL 的底层逻辑
1.1 先弄清楚它解决了什么痛点
传统实时数据处理链路的痛点,几乎人人都见过:业务方要实时指标,但离线数据仓库的批任务产出滞后,于是需要引入流处理引擎。而流处理引擎里最成熟的 Flink 和 Spark Streaming,开发语言基本都是 Java/Scala 体系,调试麻烦不说,团队如果没有 Java 功底,光环境维护和代码 review 就能折腾半个月。于是大家经常陷入一个两难:想用 Python 快速迭代,又担心性能扛不住。
Pathway 就是冲着这个矛盾去的。它是一个基于 Python 的实时数据处理框架,官方定位是“一个用于实时 ETL、流处理和增量计算的高性能 Python 框架”,底层用 Rust 实现计算引擎,对外暴露的却是纯 Python API。也就是说,你用 Pandas 一样的方式写 DataFrame 逻辑,但它跑起来并不是 Pandas 那样全量加载进内存再算,而是像流式计算一样数据进来一条处理一条。
它在 GitHub 上热度高,我认为很大程度不是因为“又有个新框架”,而是它把“Python 开发体验”和“高性能流处理”这两件事第一次比较像样地结合在了一起。社区里已经有人拿它做实时特征工程、在线推荐系统的特征拼接、时序数据的实时清洗,也有团队在做中小规模数据的流式数仓同步。
1.2 技术架构:Rust 内核 + Python 接口,增量计算模型
Pathway 性能的秘密,不在于 Python 本身变快了,而在于它把计算核心下沉到了 Rust。你在 Python 层写的select、filter、groupby,最终会被翻译成一张数据流图,由 Rust 执行引擎调度执行。这张流图是动态的,数据源有新数据到达,引擎只计算受影响的那一部分结果,而不是把全量数据重新算一遍。
这里最关键的概念是“增量计算”。拿一个很常见的场景举例:你要统计某个商品每分钟的销售总额。如果用 Spark 批处理,就是每个新批次来临时全部重算;如果用 Flink,虽然也是流式的,但在某些高阶操作里,状态管理和窗口计算仍然需要消耗大量资源去维护临时状态。Pathway 的做法更像电子表格:某个单元格的数据变了,只有依赖它的单元格会联动更新。所以对于大量“新增后追加聚合”的 ETL 场景,吞吐量天然就有优势。
Python 侧 API 的设计也很贴近日常习惯。你不需要像 Flink 那样写一堆DataStream、KeyedStream、ProcessFunction之类的抽象,直接定义输入 Schema、用pw.reducers.sum()做聚合、用pw.io.kafka.read()读数据源,看起来跟 Pandas 操作 DataFrame 几乎没有距离。但它有明确定时触发机制,数据是持续涌入而不是一次性返回结果集。
1.3 与 Pandas、Flink、Spark 的定位差异
很多人第一次看到 Pathway 时会问:这跟 Pandas 有什么区别,跟 Flink 和 Spark 又是什么关系?我用一句话概括:Pandas 是单机离线分析工具,Flink/Spark 是企业级分布式计算平台,而 Pathway 是一个“面向数据管道的实时计算引擎”。
Pandas 是一次性加载、静态计算,适合探索性数据分析和脚本化处理。Pathway 则要求你一开始就定义 Schema 和数据源,然后以流式方式持续运行。所以它承接的更多是“生产环境里的数据任务”,比如每小时同步一次数据库、每几秒消费一批 Kafka 消息做统计。
Flink 和 Spark 当然很强大,但它们的部署和运维成本摆在那里。如果你只是想快速搭一条几分钟延迟的实时 ETL 管道,数据量又没到每天几 PB 的级别,用 Flink 其实是杀鸡用牛刀。Pathway 在这种“中等规模实时数据处理”区间里,可以用 Python 一门语言把活干完,同时保持比 Mock 测试高得多的真实处理性能。
2. “吊打 Flink/Spark”的说法可信吗:性能优势与边界
2.1 基准测试看什么:端到端延迟、吞吐、资源占用
标题说“性能吊打 Flink、Spark”,这是典型的社区化表达,不要直接脑补成所有场景下 Pathway 都碾压它们。想要理性理解这句话,得看 benchmark 到底在测什么。Pathway 官方和社区常见对比指标是:端到端处理延迟、单位时间吞吐量、以及相同负载下的 CPU/内存占用。
端到端延迟指的是从数据源产生一条数据,到这条数据出现在结果表或下游系统里的时间差。Flink 做实时流处理本身延迟就低,但如果你前面链路上走了 Kafka 分区再经过连接器,每个环节都有开销。Pathway 的 Rust 引擎在数据解析、序列化和算子计算上能省掉不少时间,所以在纯计算逻辑不算复杂的情况下,它的端到端延迟可以做到很低。
吞吐量方面,Pathway 的优势更容易体现。因为它做增量计算,很多中间结果不需要反复重算,同样规格的机器,在“数据不断追加、按窗口聚合”这类典型 ETL 负载下,单位时间能吃掉的数据量确实会更优。我看到过第三方对比里单机版 Pathway 在某些场景下的吞吐表现不低于一个小型 Flink 集群,这个结论虽然不代表所有场景,但已经足够说明它单机版的实力。
资源占用是最直观的。Java 系引擎的 JVM 本身就要占不少内存,GC 调优还不能不做。Pathway 跑在 Rust 运行时上,没有 JVM 那套负担,起一个小型任务的内存峰值往往比 Flink 低得多。对于中小团队来说,这直接决定了能不能用“一台 16G 内存的机器”搞定过去需要三台机器才能扛住的实时任务。
2.2 增量计算 vs 微批处理:为什么更省
Flink 和 Spark Streaming 的底层处理模型不太一样。Spark Streaming 早年是微批处理,后来引入 Structured Streaming 后能在微批和连续处理之间切换,但常态实现仍然偏向微批。Flink 是真正的流式处理,但窗口、状态、检查点机制都比较重。
Pathway 的增量计算模型,从一开始就不是“分批处理”,而是“依赖追踪”。你可以把它想象成一张实时更新的 Excel 表:公式定义了结果怎么算,源数据单元格变化时,只重算受影响的结果单元格。这种模型在处理带累积性质的聚合、滑动窗口、去重这类任务时,省掉的计算量非常可观。
当然这种模型也有代价。为了追踪依赖关系,Pathway 需要在内存里维护数据血缘和中间状态,如果单个任务里状态无限增长,一样会面临内存压力。但它提供了外部状态存储的选项,可以把部分状态落盘,从而在成本和性能之间做权衡。
2.3 别被标题带偏:性能对比的适用场景
我必须泼一盆冷水:如果你拿 Pathway 去跑“大规模复杂图计算”或者“海量数据多阶段 join 后还要做复杂机器学习训练”,它目前肯定替代不了 Spark。Pathway 擅长的区间是“数据管道型的实时处理”:数据源相对集中,处理逻辑以过滤、转换、聚合、关联为主,对实时性要求高,但不需要跑几千个节点的大规模分布式作业。
还有就是生态问题。Flink 有大量企业级配套,包括完整的检查点机制、故障恢复、UI 监控、连接器矩阵、SQL 引擎等。Pathway 也有这些能力,但成熟度还在追赶。我的结论是:性能对比要说“在某一类实时 ETL 场景下,Pathway 可以做到优于 Flink/Spark”,这个说法更准确,也更能反映实际使用感受。
3. 30 分钟跑通一个实时 ETL Demo:订单流聚合与落库
3.1 环境准备:Python、依赖和 Kafka 实例
我用的环境是一台 Ubuntu 22.04 的 8C16G 云主机,Python 3.10,直接通过pip install pathway安装。需要注意 Pathway 的版本迭代比较快,建议固定一个大版本使用,避免 API 出现破坏性变更。我的示例环境里还装了一个单机 Kafka,使用官方 Docker 镜像跑bitnami/kafka或者wurstmeister/kafka都可以。
如果你本地还没有 Kafka 实例,可以使用下面的 Docker Compose 快速起一个:
version: '3' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:latest ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092这个配置对于本地测试足够了。如果你不想折腾 Kafka,也可以先用 CSV 文件做数据源,后面我再提怎么切换。
3.2 核心代码:定义 Schema、读取 Kafka、窗口聚合、PostgreSQL 输出
这个 Demo 的场景是:假设有一个订单事件流,业务方需要实时统计每个用户每分钟的订单金额总额,并把结果落库。我用 Pathway 消费 Kafka 里的 JSON 订单消息,做时间窗口聚合,最后写入 PostgreSQL。如果 PostgreSQL 表结构还不存在,第一次运行前可以先手动建表。
代码大致长这样:
import pathway as pw import time class OrderSchema(pw.Schema): event_id: str user_id: str product_id: str amount: float event_time: str orders = pw.io.kafka.read( rdkafka_settings={"bootstrap.servers": "localhost:9092"}, topic="orders", schema=OrderSchema, format="json", autocommit_duration_ms=1000, ) orders = orders.select( user_id=orders.user_id, amount=orders.amount, event_time=pw.this.event_time.dt.strptime("%Y-%m-%d %H:%M:%S"), window_start=pw.this.event_time.dt.timestamp() ) result = orders.windowby( pw.this.event_time, window=pw.temporal.tumbling(duration=pw.temporal.duration(seconds=60)), ).reduce( user_id=orders.user_id, total_amount=pw.reducers.sum(orders.amount), order_count=pw.reducers.count(), ) pw.io.postgres.write( result, host="localhost", port=5432, database="testdb", user="postgres", password="postgres", table="user_orders_1m", ) pw.run()这里有几个值得解释的地方。第一,pw.io.kafka.read需要传入rdkafka_settings,里面的bootstrap.servers必须和你的 Kafka 配置保持一致。第二,Schema 字段类型一定要写对,尤其是event_time这类字符串时间,得通过.dt.strptime()显式转换成时间类型,否则窗口计算会出问题。第三,窗口聚合用的是windowby加tumbling,滚动窗口,固定 60 秒一个窗口,你可以改成sliding做滑动窗口。
我实际跑下来,从启动脚本到消费到第一批结果,大约几秒钟,日志里能看到窗口聚合结果的输出。如果想直接打印结果到控制台做演示,可以换成pw.io.csv.write(result, "output.csv"),这样更容易观察结果变化。
3.3 运行与验证
在终端里先用 Python 脚本启动一个简单的 Kafka producer,每秒往orders里塞几条模拟订单,然后再启动上面的 Pathway 任务,观察结果输出:
python order_producer.py & python pathway_demo.pyPathway 任务启动后进程会一直挂着,一旦 Kafka 里有新消息,它就会增量处理,不需要手动调度。你可以同时跑多个查询,比如再做另一个窗口统计不同商品的热度,只要在同一个pw.run()之前再定义一张表就行。
这种“定义数据流图,然后一直跑”的模式,刚上手时可能会不太习惯,但一旦接受这个设定,它比你自己写循环去轮询 Kafka 或者定期跑 SQL 要省心得多。
4. 实际使用中的坑:热更新、状态恢复与生态缺失
4.1 热重载与开发体验
Pathway 给人的第一印象很好,但用到第三天你就会发现一个比较难受的地方:改代码后的热更新支持远没有脚本派来得方便。你用普通 Python 脚本时,改一个文件再跑一次就好了;Pathway 是长驻进程方式,要改动一个算子或者字段逻辑,往往需要重启整个任务,让它从最近的外部状态里恢复。
好的一点是,Pathway 的内存量不是唯一的真相,数据源和状态可以持久化到外部存储,所以重启后能接着消费 Kafka 消息继续算。但如果你没有配置状态持久化,任务一停,已经算过的窗口结果就没有了。我的建议是:开发阶段尽量用带backfilling或历史文件的方式快速验证逻辑,生产再换成长时间运行的模式。
4.2 状态管理和应用重启
状态管理是流处理里的老话题,Pathway 也避免不了。当你的应用跑了一段时间后,内存里会积累大量窗口状态、聚合中间态。如果任务重启,它会尝试从状态后端恢复。你可以在配置里指定外部状态存储,比如 SQLite 或 PostgreSQL,否则默认使用内存态。
我第一次跑一个长时间任务时,没配置外部状态存储,结果半夜机子重启,第二天发现历史窗口数据全丢了,只有重新上线后的新数据。这个问题在 Flink 里因为检查点机制成熟,大家不太容易踩到。所以如果你准备在生产环境用 Pathway,第一件事就是做好状态持久化设计和定期备份。
4.3 连接器与监控生态还不成熟
Pathway 的连接器已经涵盖了 Kafka、PostgreSQL、SQLite、CSV、Parquet、Delta Lake、S3 等常用组件,日常使用基本够。但比起 Flink 连接器社区那种“什么东西都能找到绝对兼容的 connector”的程度,Pathway 还差得远。比如某些云数据库的专有连接器、某些消息队列的认证方式,可能都需要你自己封装一层。
监控也是个短板。Flink 有丰富的 Metrics 面板、任务事件日志和报警集成,你可以很直观地看到 backpressure、checkpoint 耗时等指标。Pathway 的监控目前更多依赖日志和自定义指标导出,如果你需要精细的可观测性,得自己想办法把内部的 Kafka lag、处理延迟、错误消息数量等指标暴露出来。
5. 选型建议:什么时候该上 Pathway
5.1 与 Flink/Spark 的互补关系
不要用“谁替代谁”的眼光去看 Pathway。更准确的理解是:Pathway 填补了 Python 生态里“实时数据管道”的空缺。如果你的团队已经是 Python 技术栈,简单的实时特征、实时报表、业务指标聚合这类需求过去要硬着头皮搭 Flink,现在可以先用 Pathway 快速上线。如果数据量级真的大到需要横向扩展几十个 worker,那时再考虑引入 Spark 或 Flink 也不迟。
从资源回报率看,Pathway 特别适合“千亿级/天以下,但实时性要求在秒级到分钟级”的管道任务。它部署简单、代码量少、调试直观,能显著缩短项目周期。
5.2 该看的数据和后续发展
如果你对这个项目感兴趣,建议先看它的官方 benchmark 和文档,别只看标题。GitHub 上 API 变更也比较频繁,新版本会推出一些更贴近生产环境的功能,比如更完善的持久化状态、更多连接器、更好的监控支持。现在社区规模还不算特别大,但增长速度很快。
如果你正好在搞“实时数据入湖”、在线特征计算、或者想用 Python 替代一部分 Spark 离线 ETL,完全可以在两周内做个 PoC 试试效果。我个人的判断是,这个方向未来会和 Flink/Spark 形成互补局面,而不是谁把谁“吊打”掉。
最后再分享一个实际使用的小技巧:如果你刚开始跑 Pathway,建议先用 CSV 文件作为模拟数据源,把mode="streaming"开起来,然后往 CSV 里不断追加内容,观察结果表的变化。这个方式不需要 Kafka,也能非常直观地感受“增量计算”到底是怎么回事。等理解了它的运行模型,再换真实消息队列,会顺手很多。