十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Pathway:用Python实现高性能实时ETL,增量计算挑战Flink/Spark

Pathway:用Python实现高性能实时ETL,增量计算挑战Flink/Spark 最近这几周 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}, topicorders, schemaOrderSchema, formatjson, autocommit_duration_ms1000, ) orders orders.select( user_idorders.user_id, amountorders.amount, event_timepw.this.event_time.dt.strptime(%Y-%m-%d %H:%M:%S), window_startpw.this.event_time.dt.timestamp() ) result orders.windowby( pw.this.event_time, windowpw.temporal.tumbling(durationpw.temporal.duration(seconds60)), ).reduce( user_idorders.user_id, total_amountpw.reducers.sum(orders.amount), order_countpw.reducers.count(), ) pw.io.postgres.write( result, hostlocalhost, port5432, databasetestdb, userpostgres, passwordpostgres, tableuser_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. 选型建议什么时候该上 Pathway5.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 文件作为模拟数据源把modestreaming开起来然后往 CSV 里不断追加内容观察结果表的变化。这个方式不需要 Kafka也能非常直观地感受“增量计算”到底是怎么回事。等理解了它的运行模型再换真实消息队列会顺手很多。
返回列表