基于 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_settings=settings) 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_datetime=int(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' SECOND |
| Upsert(主键) | 发布之后若又有迟到事件到达,自动纠正结果 | 源表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 场景一:迟到但仍在耐心范围内——两个事件都被计数
假设两个发生在东村(PU=79)的上车事件,使用 10 秒窗口 + 5 秒 Watermark。事件 A 准时到达,事件 B 迟到 8 秒(乘客手机在隧道里断了信号):
事件 B 虽迟到 8 秒,但仍落在 Flink 的耐心窗口内。此时 Flink 尚未发布结果,B 被正常并入窗口,最终一次INSERT写入trips=2。
5.2 场景二:迟到超过耐心——Upsert 纠正已发布的结果
如果事件 B 迟到了 20 秒——即在 Flink 已经发布窗口结果之后才到达呢?
Flink 已经发布了trips=1,但当事件 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 ENFORCED:NOT 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),仅供参考