Data Engineering Zoomcamp Flink 的 tumbling、sliding 与 session 窗口怎么选
【免费下载链接】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 第 7 模块(Streaming)中,你会用 PyFlink 对 NYC 出租车行程流做窗口聚合。写窗口查询时第一个要回答的问题是:选 tumbling、sliding 还是 session?三种窗口把同一份数据切成不同形状的桶,结果会完全不同。这篇文章先给出课程文档中的选型依据,然后照着课程的 Redpanda + Flink + PostgreSQL 环境,实际跑通 tumbling 与 session 两个窗口作业并核对结果,最后给出 sliding 窗口(HOP)的 SQL 模板和适用场景。
选型依据:三种窗口的区别与适用场景
课程文档 Understanding window types 对三种窗口的定义是:窗口类型决定一个事件属于一个固定桶、多个重叠桶,还是一段以不活动边界限定的突发(burst)。
| 窗口类型 | 大小 | 事件归属 | 文档给出的使用场景 |
|---|---|---|---|
| Tumbling(翻滚) | 固定、不重叠 | 每个事件恰好属于一个窗口 | 每小时行程计数、每日收入汇总 |
| Sliding(滑动) | 固定、重叠 | 一个事件可以属于多个窗口 | 找峰值谷值("任意 1 小时窗口内的流量峰值?")、移动平均、 surge 检测(如网约车动态加价) |
| Session(会话) | 不固定,由不活动间隙决定 | 事件按不活动时长归入同一会话 | 把用户行为聚成会话(sessionization),用于行为分析 |
按这个对照表选择:需要按固定时间点(每小时、每 5 分钟)出汇总,用 tumbling;需要"任意 N 分钟/小时窗口"的重叠视角来找极值或做移动平均,用 sliding;聚合单位不是固定时刻,而是"一段连续活动、被静默隔开",用 session。
下面用课程 2026 队列的作业环境实际跑通 tumbling 和 session。
准备:启动 Redpanda + Flink + PostgreSQL 环境
前置条件来自 workshop README:
- Docker 和 Docker Compose
- uv(用于运行 Python 依赖)
- 一个 SQL 客户端:pgcli(
uvx pgcli)、DBeaver、pgAdmin 或 DataGrip
作业说明给出的启动方式:
cd 07-streaming/workshop/ docker compose build docker compose up -d启动后得到四个服务:Redpanda(Kafka 兼容 broker,localhost:9092)、Flink Job Manager(http://localhost:8081)、Flink Task Manager、PostgreSQL(localhost:5432,用户postgres,密码postgres)。
两个注意事项:
- 容器名(如
workshop-redpanda-1)依赖目录叫workshop,如果你改过目录名,命令里的容器名要相应调整。 - 之前跑过作业、留有旧容器或数据卷时,先做干净启动。
docker compose down -v会删除数据卷(包括容器里的数据),再重新build和up -d:
docker compose down -v docker compose build docker compose up -d用docker compose ps确认四个服务均为Up。
发送 green taxi 数据到green-tripstopic
先创建 topic:
docker exec -it workshop-redpanda-1 rpk topic create green-trips然后运行 producer。可以直接使用课程提供的现成脚本 producer.py:它读取 green_tripdata_2025-10.parquet(49,416 行,脚本内置数据 URL),只保留lpep_pickup_datetime、lpep_dropoff_datetime、PULocationID、DOLocationID、passenger_count、trip_distance、tip_amount、total_amount共 8 列,把 datetime 转成字符串后逐行以 JSON 发到green-trips。在07-streaming/workshop/项目目录内(uv sync之后)执行:
python producer.py文档给出的参考结果是Sent 49416 messages,耗时约 10 秒(随机器不同会变化)。
如果 topic 被发送过多次,数据会重复。文档给出的处理方式是删除并重建 topic:docker exec -it workshop-redpanda-1 rpk topic delete green-trips后重新执行上一条 create 命令。注意rpk topic delete会删掉该 topic 及其中的数据。
Tumbling 窗口实战:5 分钟窗口统计各上车点行程数
"哪个上车点在单个 5 分钟窗口内行程最多"是典型的固定时间点、不重叠问题,对应 tumbling。
先在 PostgreSQL 建结果表(在 pgcli 或docker compose exec postgres psql -U postgres -d postgres里执行):
CREATE TABLE tumbling_pickup_counts ( window_start TIMESTAMP(3), PULocationID INT, num_trips BIGINT, PRIMARY KEY (window_start, PULocationID) NOT ENFORCED );把 tumbling_job.py 复制到07-streaming/workshop/src/job/目录(该目录挂载到 Flink 容器的/opt/src/job/)。作业的关键部分:
- 源表 DDL 中把字符串时间转成事件时间并定义 watermark:
event_timestamp AS TO_TIMESTAMP(lpep_pickup_datetime, 'yyyy-MM-dd HH:mm:ss'), WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL '5' SECOND- 窗口函数为 5 分钟翻滚窗口:
FROM TABLE( TUMBLE(TABLE green_trips, DESCRIPTOR(event_timestamp), INTERVAL '5' MINUTE) ) GROUP BY window_start, PULocationIDenv.set_parallelism(1)——因为green-tripstopic 只有 1 个分区,更高的并行度会留下空闲的 subtask,导致 watermark 无法推进、窗口不出结果。'scan.startup.mode' = 'earliest-offset',从头读完 topic 里的全部数据。
提交作业:
docker exec -it workshop-jobmanager-1 \ flink run -py /opt/src/job/tumbling_job.pyFlink 流作业是持续运行的:让作业跑一两分钟,直到 PostgreSQL 出现结果,然后到 Flink UI(http://localhost:8081)取消作业。
验证查询:
SELECT PULocationID, num_trips FROM tumbling_pickup_counts ORDER BY num_trips DESC LIMIT 3;solutions.md 记录的预期答案:PULocationID 74,在最忙的 5 分钟窗口内有 15 趟行程。
关于 watermark,Aggregation with tumbling windows 的解释是:窗口定义"数什么",watermark 定义"何时发布结果"。上面的 5 秒 tolerance 是等待迟到事件的耐心值;watermark 越大对迟到事件越宽容,但看到结果的等待也更长。文档称 5 秒是合理的默认值,生产环境应根据数据实际的乱序程度调整。
Session 窗口实战:按 5 分钟不活动间隙聚成会话
"连续活动的突发"不是固定时间点,而是被静默隔开的事件序列——这正是 session 窗口的场景:按PULocationID分组,某上车点的行程只要相邻间隔不超过 5 分钟就归入同一会话,间隔超过 5 分钟会话关闭。
先建结果表:
CREATE TABLE session_pickup_counts ( session_start TIMESTAMP(3), session_end TIMESTAMP(3), PULocationID INT, num_trips BIGINT, PRIMARY KEY (session_start, PULocationID) NOT ENFORCED );把 session_job.py 复制到07-streaming/workshop/src/job/,窗口函数是:
FROM TABLE( SESSION(TABLE green_trips_session PARTITION BY PULocationID, DESCRIPTOR(event_timestamp), INTERVAL '5' MINUTE) ) GROUP BY window_start, window_end, PULocationID与 tumbling 作业相比,SESSION多了一个PARTITION BY PULocationID(每个上车点独立划分会话),并且输出带session_start/session_end两列。同样地,源表定义 5 秒的WATERMARK,且env.set_parallelism(1)是必须的——solutions.md 特别指出:green-trips只有 1 个分区,并行度高于 1 时空闲 subtask 会阻止 watermark 推进,session 窗口根本不会出结果。
提交:
docker exec -it workshop-jobmanager-1 \ flink run -py /opt/src/job/session_job.py等一两分钟后到 Flink UI 取消作业,然后查询:
SELECT PULocationID, num_trips, session_start, session_end FROM session_pickup_counts ORDER BY num_trips DESC LIMIT 3;文档记录的预期答案:最长会话为 81 趟(PULocationID 74,2025-10-08 上午)。
把两个结果放在一起看,选型差异就很直观:同一份 2025 年 10 月 green taxi 数据,5 分钟固定窗口(tumbling)下最忙的单窗口只有 15 趟,而 5 分钟不活动间隙的会话(session)把连续活动的行程串起来,最长会话达到 81 趟。如果你的问题是"任意固定时间点的量",用前者;如果你的问题是"一次连续行为有多长",用后者。
Sliding 窗口:何时用 HOP,怎么写
"任意 1 小时窗口内的峰值流量"这类问题,tumbling 回答不了:00:00–01:00 是一个 1 小时窗口,00:15–01:15 也是,重叠的滑动窗口能同时表达所有这些起点。文档给出的 SQL 模板是HOP:
HOP(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL '15' MINUTE, INTERVAL '1' HOUR)即 1 小时的窗口、每 15 分钟滑动一次;重叠意味着一个事件会落入多个窗口。适用场景是找峰值谷值、移动平均和 surge 检测(如网约车动态加价)。
需要说明:仓库没有提供现成的 sliding 作业文件,只有这段 SQL 模板。要实际运行,可以参照 tumbling_job.py 的结构,把TUMBLE(...)替换为HOP(...),源表仍需保留event_timestamp计算列和WATERMARK定义,sink 表的主键需要能容纳重叠窗口的多行结果(不同窗口起点)。
常见问题与限制
- 窗口不出结果:先检查并行度与 topic 分区。
green-trips是 1 个分区,作业必须env.set_parallelism(1);否则空闲 subtask 让 watermark 停住,窗口永远不会触发(对 tumbling 和 session 都适用)。 - 结果重复:作业使用
earliest-offset,每次提交都从头读 topic。如果 producer 跑了多次,先删除并重建 topic(rpk topic delete green-trips),这会删掉 topic 数据。 - 作业是常驻的:Flink 流作业不会自动结束,跑完验证后到 Flink UI(http://localhost:8081)取消。
- watermark 是折中项:5 秒 tolerance 等待的是几秒内的迟到事件;更大意味着更完整但结果发布更慢。文档只给出"5 秒是合理默认值",没有给出针对更大延迟场景的调优参数。
- sliding 无现成作业:仓库只有
HOP的 SQL 模板,需要自行参照 tumbling/session 作业组装,本文只给出替换方式,不提供完整作业文件。
完成验证后,可以按课程模块的 13-cleanup.md 停掉容器。回头看选型本身:固定时间点归一桶用 tumbling,重叠视角找极值用 sliding,按不活动划分行为段用 session——三个问题、三种窗口,本文的 tumbling 与 session 路径都基于同一份 2025 年 10 月 green taxi 数据集和同一套 workshop 基础设施,结果差异完全来自窗口形状,可以直接对照复现。
【免费下载链接】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),仅供参考