ARTICLE DETAIL

资讯详情

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

Data Engineering Zoomcamp Flink 的 tumbling、sliding 与 session 窗口怎么选

Data Engineering Zoomcamp Flink 的 tumbling、sliding 与 session 窗口怎么选 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 READMEDocker 和 Docker Composeuv用于运行 Python 依赖一个 SQL 客户端pgcliuvx pgcli、DBeaver、pgAdmin 或 DataGrip作业说明给出的启动方式cd 07-streaming/workshop/ docker compose build docker compose up -d启动后得到四个服务RedpandaKafka 兼容 brokerlocalhost:9092、Flink Job Managerhttp://localhost:8081、Flink Task Manager、PostgreSQLlocalhost:5432用户postgres密码postgres。两个注意事项容器名如workshop-redpanda-1依赖目录叫workshop如果你改过目录名命令里的容器名要相应调整。之前跑过作业、留有旧容器或数据卷时先做干净启动。docker compose down -v会删除数据卷包括容器里的数据再重新build和up -ddocker compose down -v docker compose build docker compose up -d用docker compose ps确认四个服务均为Up。发送 green taxi 数据到green-tripstopic先创建 topicdocker exec -it workshop-redpanda-1 rpk topic create green-trips然后运行 producer。可以直接使用课程提供的现成脚本 producer.py它读取 green_tripdata_2025-10.parquet49,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 被发送过多次数据会重复。文档给出的处理方式是删除并重建 topicdocker 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 中把字符串时间转成事件时间并定义 watermarkevent_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 UIhttp://localhost:8081取消作业。验证查询SELECT PULocationID, num_trips FROM tumbling_pickup_counts ORDER BY num_trips DESC LIMIT 3;solutions.md 记录的预期答案PULocationID 74在最忙的 5 分钟窗口内有 15 趟行程。关于 watermarkAggregation 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 742025-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 模板是HOPHOP(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 跑了多次先删除并重建 topicrpk topic delete green-trips这会删掉 topic 数据。作业是常驻的Flink 流作业不会自动结束跑完验证后到 Flink UIhttp://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),仅供参考
返回列表