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

资讯详情

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

基于 PyFlink 的 Tumbling Window 聚合:从窗口语义到 Watermark 与 Upsert 的完整实战

基于 PyFlink 的 Tumbling Window 聚合:从窗口语义到 Watermark 与 Upsert 的完整实战 基于 PyFlink 的 Tumbling Window 聚合从窗口语义到 Watermark 与 Upsert 的完整实战【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp本篇技术指南聚焦于 Data Engineering Zoomcamp 流处理模块中的核心实战课——使用 PyFlink SQL 对 Kafka/Redpanda 中的出租车行程事件做按小时、按上车地点的滚动窗口Tumbling Window聚合并将结果以 Upsert 语义写入 PostgreSQL。你将理解窗口、Watermark、主键 Upsert 三者如何协同工作掌握从建表、提交作业到验证结果的完整操作流程并能在纯 Python 消费者无法胜任的窗口聚合场景中直接复用这套方案。一、为什么需要窗口聚合前面章节中我们先后用纯 Python 消费者04-consume-messages-with-python.md和 Flink 透传作业08-the-pass-through-flink-job.md实现了读 Kafka → 写 PostgreSQL的数据搬运。这类逐条处理pass-through不需要记忆任何历史状态——来一条、写一条。但本次任务发生了质变我们要统计每个上车地点PULocationID每小时产生了多少趟出租车行程、合计多少营收。这要求系统必须把事件按时间切分到固定的桶window中在桶内维护计数和求和状态在合适的时机把结果发布出来处理迟到事件对已发布结果的修正。用纯 Python 消费者实现这些需要自己维护窗口状态、处理乱序与迟到、管理失败恢复还要手写 Upsert SQL——正如原文档所说这几乎是重写一套流处理框架。而用 Flink这只是一条 SQL 查询。二、准备阶段取消旧作业并创建 PostgreSQL 聚合表在编写聚合作业之前先取消任何正在运行的旧作业透传作业会持续占用资源并写processed_events表。然后在 PostgreSQL 中创建聚合结果表CREATE TABLE processed_events_aggregated ( window_start TIMESTAMP, PULocationID INTEGER, num_trips BIGINT, total_revenue DOUBLE PRECISION, PRIMARY KEY (window_start, PULocationID) );这张表有两个至关重要的设计决策直接决定了整个管线的正确性PULocationID被纳入表结构因为聚合同时按时间窗口和上车地点分组GROUP BY window_start, PULocationID所以分组键必须全部出现在输出表和主键中。复合主键(window_start, PULocationID)这是启用 Upsert 行为的关键。当 Flink 针对同一个窗口发出更新后的计数时PostgreSQL 会更新已有行而不是插入重复行。这一点非常重要因为迟到事件可能导致 Flink 重新评估一个已经发布过结果的窗口——有了主键 Upsert修正后的计数会自动替换旧值无需人工干预。仓库中的完整作业代码位于 aggregation_job.py下文将逐段拆解。三、聚合作业全貌一段完成窗口聚合的 Flink SQL创建src/job/aggregation_job.py完整代码与仓库 aggregation_job.py 保持一致from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import EnvironmentSettings, StreamTableEnvironment def create_events_source_kafka(t_env): table_name events source_ddl f CREATE TABLE {table_name} ( PULocationID INTEGER, DOLocationID INTEGER, trip_distance DOUBLE, total_amount DOUBLE, tpep_pickup_datetime BIGINT, event_timestamp AS TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3), WATERMARK for event_timestamp as event_timestamp - INTERVAL 5 SECOND ) WITH ( connector kafka, properties.bootstrap.servers redpanda:29092, topic rides, scan.startup.mode earliest-offset, properties.auto.offset.reset earliest, format json ); t_env.execute_sql(source_ddl) return table_name def create_events_aggregated_sink(t_env): table_name processed_events_aggregated sink_ddl f CREATE TABLE {table_name} ( window_start TIMESTAMP(3), PULocationID INT, num_trips BIGINT, total_revenue DOUBLE, PRIMARY KEY (window_start, PULocationID) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:postgresql://postgres:5432/postgres, table-name {table_name}, username postgres, password postgres, driver org.postgresql.Driver ); t_env.execute_sql(sink_ddl) return table_name def log_aggregation(): env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) env.set_parallelism(3) settings EnvironmentSettings.new_instance().in_streaming_mode().build() t_env StreamTableEnvironment.create(env, environment_settingssettings) try: source_table create_events_source_kafka(t_env) aggregated_table create_events_aggregated_sink(t_env) t_env.execute_sql(f INSERT INTO {aggregated_table} SELECT window_start, PULocationID, COUNT(*) AS num_trips, SUM(total_amount) AS total_revenue FROM TABLE( TUMBLE(TABLE {source_table}, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR) ) GROUP BY window_start, PULocationID; ).wait() except Exception as e: print(Writing records from Kafka to JDBC failed:, str(e)) if __name__ __main__: log_aggregation()3.1 Kafka 源表新增的两行事件时间与 Watermark与透传作业 pass_through_job.py 中的 Kafka 源表相比聚合作业多了两行关键声明event_timestamp AS TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3), WATERMARK for event_timestamp as event_timestamp - INTERVAL 5 SECONDevent_timestamp AS TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3)一个计算列computed column把生产者发送的自纪元起的毫秒时间戳tpep_pickup_datetime BIGINT来自 models.py 中Ride数据类生产者把它编码为 epoch 毫秒转换为 Flink 的时间戳类型。参数3表示毫秒精度TIMESTAMP(3)。WATERMARK for event_timestamp as event_timestamp - INTERVAL 5 SECOND定义事件时间列上的水位线Watermark它告诉 Flink何时可以发布窗口结果。注意与透传作业的另两处差异scan.startup.mode earliest-offset配合properties.auto.offset.reset earliest从 Kafka 主题最旧的偏移量开始消费确保rides主题中已有的历史数据1000 条出租车记录全部被读取并参与聚合。透传作业使用latest-offset只消费新消息这一差异的详细讨论见 09-offsets-earliest-vs-latest.md。env.set_parallelism(3)以 3 个并行副本处理数据配合 docker-compose.yml 中 TaskManager 的taskmanager.numberOfTaskSlots: 15和parallelism.default: 3配置。3.2 事件时间戳从哪来聚合必须基于事件发生的时间事件时间而不是 Flink 收到消息的机器时间。在 models.py 中生产者把行程记录的时间字段转换为 epoch 毫秒tpep_pickup_datetimeint(row[tpep_pickup_datetime].timestamp() * 1000),而 producer.py 从 NYC 出租车公开数据集yellow_tripdata_2025-11.parquet取前 1000 行读取数据并序列化为 JSON 发送到rides主题。TO_TIMESTAMP_LTZ负责把这个毫秒整数还原为带时区的可比较时间戳从而作为窗口切分的依据。四、Watermark流处理中何时发布的触发器窗口Window定义了统计什么——一个 1 小时的出租车行程桶。但在流式场景中事件是持续到达的Flink 怎么知道14:00–15:00 这个小时的计数该发布了它不能只看系统时钟因为有些事件会迟到。如果没有触发器Flink 会无限期累积数据永远不会向 PostgreSQL 写入任何结果。Watermark 就是那个触发器。在 SQL 中WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 5 SECOND ^^^^^^^^^^^^^^^^^^^ patience 5 secondsWatermark 始终比 Flink 见过的最新事件时间戳落后 5 秒。当 Watermark 越过某个窗口的结束边界时Flink 就发布该窗口的结果。这 5 秒是给掉队者的耐心——那些事件时间在窗口结束之前、但到达时间晚了若干秒的事件。三个组件各司其职、协同工作组件职责实现位置窗口Window定义往哪个桶里计数1 小时TUMBLE(..., INTERVAL 1 HOUR)Watermark定义何时发布结果触发器WATERMARK FOR event_timestamp AS ... - INTERVAL 5 SECONDUpsert主键发布之后若又有迟到事件到达自动纠正结果源表PRIMARY KEY 目标表主键4.1 窗口与 Watermark 的 SQL 语法TUMBLE(TABLE {source_table}, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR)TUMBLE创建固定大小、不重叠的滚动窗口例如 [00:00, 01:00)、[01:00, 02:00)…DESCRIPTOR(event_timestamp)必须引用定义了 WATERMARK 的列否则 Flink 无法推进窗口发布INTERVAL 1 HOUR设定窗口大小此处为 1 小时。五、时序推演迟到事件在两种情形下的命运原文档用两个 Mermaid 时序图把窗口 Watermark Upsert 的行为讲得非常透彻这里完整保留并展开说明。5.1 场景一迟到但仍在耐心范围内——两个事件都被计数假设两个发生在东村PU79的上车事件使用 10 秒窗口 5 秒 Watermark。事件 A 准时到达事件 B 迟到 8 秒乘客手机在隧道里断了信号事件 B 虽迟到 8 秒但仍落在 Flink 的耐心窗口内。此时 Flink 尚未发布结果B 被正常并入窗口最终一次INSERT写入trips2。5.2 场景二迟到超过耐心——Upsert 纠正已发布的结果如果事件 B 迟到了 20 秒——即在 Flink 已经发布窗口结果之后才到达呢Flink 已经发布了trips1但当事件 B 终于到达时主键让 Flink 能够发送一条修正PostgreSQL 把该行从 1 更新为 2。如果没有主键即 append-only sink事件 B 会被直接丢弃——因为在追加模式下Flink 无法重新打开一个已经发布的窗口。这正是后续章节 11-late-events-and-upserts.md 要实验的现象使用实时生产者producer_realtime.py约 20% 的事件带 3–10 秒的过去时间戳持续灌数据时可以watch到旧窗口的计数随着迟到事件到达而增长——每次增长都是一次主键 Upsert。5.3 延迟与完整性的权衡Watermark 是一个延迟latency与完整性completeness之间的权衡Watermark 越大等待迟到事件的耐心越足、结果越完整但代价是你要等更久才能看到任何结果。5 秒是一个合理的默认值。在生产环境中应根据数据真实的乱序程度来调优——数据越乱序需要的 patience 越大。六、与透传作业的其他差异除 Watermark 与计算列外聚合作业还有几处值得注意Sink 带PRIMARY KEY (...)且NOT ENFORCEDNOT ENFORCED表示该主键由 Flink 声明但不强制校验其作用是在 Flink JDBC 连接器中启用 Upsert 行为。对比透传作业的 sink pass_through_job.py无主键、append-only这是两者在 sink 定义上最大的区别。earliest-offset从头读取 Kafka 中已有的全部数据。env.set_parallelism(3)3 个并行实例处理数据与 TaskManager 配置呼应。TUMBLE窗口函数产生固定大小、互不重叠的窗口DESCRIPTOR必须指向声明了 Watermark 的列INTERVAL 1 HOUR决定窗口大小。env.enable_checkpointing(10 * 1000)每 10 秒做一次状态快照。对窗口作业而言checkpoint 会把尚未关闭的窗口状态序列化到磁盘——如果作业运行 2 分钟时失败此时 5 分钟窗口还没关重启后能带着半满的窗口原地恢复而不是从头再来。七、提交作业并验证结果7.1 提交聚合作业docker compose exec jobmanager ./bin/flink run \ -py /opt/src/job/aggregation_job.py \ --pyFiles /opt/src -d其中--pyFiles /opt/src把源码目录挂入作业类路径对应 docker-compose.yml 中 JobManager 容器将./src/挂载为/opt/src-d表示后台分离运行。也可以到http://localhost:8081的 Flink Web UI 查看作业运行状态。7.2 发送数据uv run python src/producers/producer.py该生产者读取 2025 年 11 月的纽约黄色出租车数据前 1000 行序列化为 JSON 后以约 10ms 间隔逐条发送到rides主题。7.3 查询聚合结果等待约 15 秒让窗口关闭Watermark 推进越过窗口边界然后查询SELECT window_start, count(*) as locations, sum(num_trips) as total_trips, round(sum(total_revenue)::numeric, 2) as revenue FROM processed_events_aggregated GROUP BY window_start ORDER BY window_start;预期输出形如window_start | locations | total_trips | revenue ------------------------------------------------------- 2025-11-01 00:00:00 | ... 2025-11-01 01:00:00 | ... ...1000 条出租车行程被按上车地点归入 1 小时的滚动窗口中。每一行展示该小时内有行程的地点数量locations、总行程数total_trips以及合计营收revenue。八、回顾为什么用 Flink 做这件事把同样的需求交给纯 Python 消费者你需要自己实现窗口切分与聚合逻辑按事件时间而不是到达时间乱序与迟到事件的处理策略长时间运行的窗口状态管理含故障恢复面向 PostgreSQL 的 Upsert SQL 编写。而用 PyFlink以上全部收敛为一段 SQL 查询加两张 DDL 表声明。窗口负责统计什么Watermark 负责何时发布主键 Upsert 负责发布后如何纠错——理解这三者的配合就掌握了流式窗口聚合的核心模型也就能在此基础上继续探索更多窗口类型见 12-understanding-window-types.md与更复杂的迟到事件处理策略。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表