
Data Engineering Zoomcamp 2026 模块七Redpanda PyFlink 流式作业官方题解实战指南【免费下载链接】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 2026 第 7 周流处理模块家庭作业的完整题解基于 solutions.md 展开并结合仓库中的源码与配置逐行剖析。你将掌握一条端到端的实时数据处理链路用 RedpandaKafka 协议兼容作为消息中间件、用 Python producer/consumer 完成数据的生产与消费、再借助 PyFlink 的 Table API 以 tumbling 窗口和 session 窗口完成聚合分析并写入 PostgreSQL。读完本文你能独立复现六道题的完整答案并理解并行度、水位线watermark与窗口触发机制这些流式处理中的关键细节。背景作业数据与基础设施本作业使用 2025 年 10 月的 Green Taxi Trip 数据parquet 格式完整背景见 homework.md。基础设施直接复用 workshop 目录下的 Docker 编排启动命令为cd 07-streaming/workshop/ docker compose build docker compose up -d如需彻底清理旧容器与卷可先执行docker compose down -v再重新构建。从 docker-compose.yml 可以看到整套环境包含四个服务服务镜像对外端口用途redpandaredpandadata/redpanda:v25.3.99092Kafka 协议Kafka 兼容的消息代理jobmanagerpyflink-workshop本地构建8081Flink 作业管理与 Web UItaskmanagerpyflink-workshop本地构建—Flink 任务执行15 个槽位默认并行度 3postgrespostgres:185432结果存储用户/密码均为postgres关键细节容器内 Flink 通过redpanda:29092访问 Kafka 协议端口宿主机则通过localhost:9092./src/目录被挂载到 Flink 容器的/opt/src这正是后续作业文件放置位置的由来。容器名如workshop-redpanda-1基于目录名workshop生成若重命名目录需相应调整命令。问题 1确认 Redpanda 版本进入 Redpanda 容器执行rpkRedpanda 自带的管理 CLI查看版本docker exec -it workshop-redpanda-1 rpk version答案v25.3.9与 docker-compose.yml 中的镜像标签redpandadata/redpanda:v25.3.9完全一致。这个答案也提示了一个验证技巧任何容器内工具版本都可以与编排文件中的镜像 tag 交叉核对。问题 2向 Redpanda 生产数据创建 Topicdocker exec -it workshop-redpanda-1 rpk topic create green-trips生产者实现剖析仓库中的 producer.py 完整展示了数据生产流程核心要点如下读取并裁剪列用pd.read_parquet(url, columns[...])直接从云端读取 parquet仅保留作业要求的 8 个字段lpep_pickup_datetime、lpep_dropoff_datetime、PULocationID、DOLocationID、passenger_count、trip_distance、tip_amount、total_amount。时间序列化datetime 列必须先转成字符串才能 JSON 序列化使用df[lpep_pickup_datetime].dt.strftime(%Y-%m-%d %H:%M:%S)。空值处理passenger_count用fillna(0).astype(int)填充 NaN 并转为整型避免 JSON 中出现NaN导致下游解析失败。发送与计时KafkaProducer配置value_serializerjson_serializerjson.dumps(...).encode(utf-8)逐行producer.send(topic_name, valuemessage)后调用producer.flush()确保全部落盘并用time()记录耗时。运行方式python producer.py答案数据集共 49,416 行发送全程约 10 秒视机器性能略有浮动。这也印证了flush()的语义——它阻塞直到缓冲区中的全部消息发送完成是准确计时的前提。问题 3消费者统计 trip_distance 5 的行程数consumer.py 给出了一个极简但完整的消费者模板consumer KafkaConsumer( topic_name, bootstrap_servers[server], auto_offset_resetearliest, # 从头读取所有消息 group_idgreen-trips-homework, value_deserializerjson_deserializer, consumer_timeout_ms10000, # 10 秒无新消息即停止 )关键参数语义auto_offset_resetearliest从分区最旧偏移量开始消费保证统计到全部 49,416 条消息consumer_timeout_ms10000消费完数据后 10 秒内没有新消息就自动退出循环避免程序永久挂起计数器在循环内对trip[trip_distance] 5.0累加。答案trip_distance 5.0的行程共 8506 条总消息数 49,416 条。PyFlink 部分三个作业的公共骨架问题 46 的三个 Flink 作业结构高度一致理解公共骨架即可一通百通。以 tumbling_job.py 为例它由四部分组成1. 环境初始化——显式设置并行度与检查点env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) env.set_parallelism(1)2. Kafka Source DDL——把字符串时间戳转换为事件时间并声明水位线CREATE TABLE green_trips ( lpep_pickup_datetime VARCHAR, ... event_timestamp AS TO_TIMESTAMP(lpep_pickup_datetime, yyyy-MM-dd HH:mm:ss), WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 5 SECOND ) WITH ( connector kafka, properties.bootstrap.servers redpanda:29092, topic green-trips, scan.startup.mode earliest-offset, properties.auto.offset.reset earliest, format json );计算列computed columnevent_timestamp用TO_TIMESTAMP把lpep_pickup_datetime字符串按yyyy-MM-dd HH:mm:ss格式解析成时间戳WATERMARK FOR ... - INTERVAL 5 SECOND声明 5 秒水位线容忍度允许迟到 5 秒以内的数据仍能参与窗口计算。3. JDBC Sink DDL——connector jdbc连接jdbc:postgresql://postgres:5432/postgres用户名密码均为postgres并用PRIMARY KEY (...) NOT ENFORCED声明主键以支持结果更新upsert。4. 窗口聚合 SQL——见各题详解。必须注意的两个坑并行度必须为 1green-tripstopic 只有 1 个分区。若并行度过高空闲的 consumer 子任务会阻止水位线前进导致窗口永不触发详细原因见问题 5作业是流式常驻的提交后让作业运行 12 分钟直到结果写入 PostgreSQL再从 Flink Web UIhttp://localhost:8081取消作业或直接查询结果表。若多次发送过数据需先删除重建 topic 避免重复docker exec -it workshop-redpanda-1 rpk topic delete green-trips问题 45 分钟 Tumbling 窗口统计各 PULocationID 的行程数建表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/该目录已挂载到容器/opt/src/job/然后提交docker exec -it workshop-jobmanager-1 \ flink run -py /opt/src/job/tumbling_job.py窗口聚合 SQLINSERT INTO tumbling_pickup_counts SELECT window_start, PULocationID, COUNT(*) AS num_trips FROM TABLE( TUMBLE(TABLE green_trips, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE) ) GROUP BY window_start, PULocationID;Tumbling 窗口是固定长度、互不重叠的时间桶5 分钟一个DESCRIPTOR(event_timestamp)指明按事件时间对齐窗口边界COUNT(*)统计每个窗口内每个上车地点的行程数。查询与答案SELECT PULocationID, num_trips FROM tumbling_pickup_counts ORDER BY num_trips DESC LIMIT 3;答案PULocationID 74其最繁忙的 5 分钟窗口内有 15 趟行程。问题 5Session 窗口找出最长连续会话建表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 );提交作业与聚合 SQL将 session_job.py 复制到07-streaming/workshop/src/job/后提交docker exec -it workshop-jobmanager-1 \ flink run -py /opt/src/job/session_job.pyINSERT INTO session_pickup_counts SELECT window_start AS session_start, window_end AS session_end, PULocationID, COUNT(*) AS num_trips FROM TABLE( SESSION(TABLE green_trips_session PARTITION BY PULocationID, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE) ) GROUP BY window_start, window_end, PULocationID;Session 窗口与 Tumbling 窗口的本质区别它没有固定长度而是按事件间隙切分——当两个事件之间的时间间隔超过INTERVAL 5 MINUTE5 分钟 gap时前一个会话关闭、新会话开启。PARTITION BY PULocationID表示按上车地点分别维护会话这样每个地点的活动才能独立成组。查询与答案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 上午。深入理解为什么并行度必须是 1solutions 文档特别强调了一个易踩的坑session 作业必须使用并行度 1。原因在于green-trips只有 1 个分区若并行度大于 1其余消费者子任务在分配不到数据时成为空闲子任务而 Flink 的水位线只会在所有输入都越过某阈值时才会整体推进空闲子任务无法产生新的水位线导致事件时间停滞session 窗口永远不会触发关闭。同理这也解释了为什么问题 4 与问题 6 的作业同样设置env.set_parallelism(1)。问题 61 小时 Tumbling 窗口计算每小时小费总额建表CREATE TABLE hourly_tips ( window_start TIMESTAMP(3), total_tips DOUBLE, PRIMARY KEY (window_start) NOT ENFORCED );提交作业与聚合 SQL将 tip_job.py 复制到07-streaming/workshop/src/job/后提交docker exec -it workshop-jobmanager-1 \ flink run -py /opt/src/job/tip_job.pyINSERT INTO hourly_tips SELECT window_start, SUM(tip_amount) AS total_tips FROM TABLE( TUMBLE(TABLE green_trips_tips, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR) ) GROUP BY window_start;与问题 4 相比这里把窗口粒度从 5 分钟放大到 1 小时聚合函数从COUNT(*)换成SUM(tip_amount)且不再按PULocationID分组作业要求统计全地点每小时的小费总额。查询与答案SELECT window_start, total_tips FROM hourly_tips ORDER BY total_tips DESC LIMIT 3;答案小费总额最高的小时是2025-10-16 18:00:00该小时小费总计约$524.96。总结一套可复用的流处理模式回顾全部六道题可以提炼出一条完整的实战方法论基础设施即代码用 docker compose 一键拉起 Redpanda FlinkJob/Task Manager PostgreSQL版本信息如 Redpandav25.3.9可直接从编排文件核对生产侧parquet → pandas 列裁剪 → datetime 转字符串 → JSON 序列化 → KafkaProducer 发送 flush()计时消费侧auto_offset_resetearliest全量消费consumer_timeout_ms控制退出时机Flink 侧统一采用Kafka Source DDL 事件时间计算列 5 秒水位线 JDBC Sink DDL 窗口聚合 SQL的模板只需替换窗口类型TUMBLE/SESSION、窗口长度5 分钟/1 小时与聚合函数COUNT/SUM即可复用到其他指标并行度纪律单分区 topic 配合窗口计算时必须set_parallelism(1)这是水位线能否推进、窗口能否触发的关键前提。如需查看作业源码的完整实现可进一步阅读 tumbling_job.py、session_job.py、tip_job.py 与 workshop 侧的 aggregation_job.py其中包含事件时间、水位线与窗口的对照演示以及 workshop README 了解完整环境搭建步骤。【免费下载链接】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),仅供参考