Data Engineering Zoomcamp 2027 第 7 周作业实战:用 Redpanda 与 PyFlink 构建绿色出租车实时流处理管道

发布时间:2026/9/12 7:20:03

Data Engineering Zoomcamp 2027 第 7 周作业实战:用 Redpanda 与 PyFlink 构建绿色出租车实时流处理管道 Data Engineering Zoomcamp 2027 第 7 周作业实战用 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 2027 年 cohort 第 7 模块流处理的作业展开核心任务是以 RedpandaKafka 兼容消息中间件为消息枢纽、以 PyFlink 为流处理引擎对 2025 年 10 月的纽约绿色出租车Green Taxi行程数据完成生产 → 消费 → 窗口聚合的完整实时链路演练。读完本文你将掌握 Kafka Producer/Consumer 的 Python 写法、Flink SQL 中 tumbling/session 窗口与 watermark 的实战用法并能在本地 Docker 环境中独立复现并验证全部 6 道作业题。作业概览一条完整的实时流处理链路本作业建立在 07 模块 workshop 的基础设施之上目标是复用同一套 Docker 编排把 workshop 中基于黄色出租车yellow taxi、tpep_前缀字段、ridestopic的示例迁移到绿色出租车green taxi、lpep_前缀字段、green-tripstopic数据上从而检验三个核心技能点消息生产与消费用kafka-python库把 parquet 数据写入 Redpanda topic再以消费者身份读回并做简单过滤统计事件时间与窗口在 PyFlink 中以字符串形式的时间戳声明事件时间用 5 秒 watermark 容忍乱序分别实现 tumbling滚动与 session会话窗口聚合结果落库通过 Flink JDBC 连接器把聚合结果写入 PostgreSQL 并查询验证。数据来源为纽约 TLC 公开的绿色出租车行程快照文件green_tripdata_2025-10.parquet2025 年 10 月整月数据作业环境由 07-streaming/workshop 目录下的 Docker Compose 提供作业中引用的完整代码镜像位于 cohorts/2027/07-streaming/code。环境准备复用 workshop 的流处理基础设施作业直接复用 workshop 的环境无需从零搭建。在仓库根目录下进入 workshop 目录并启动全部服务cd 07-streaming/workshop/ docker compose build docker compose up -dbuild会基于 Dockerfile.flink 构建带 Python 与 PyFlink 的自定义 Flink 镜像起点是官方flink:2.2.0-scala_2.12-java17内置 Python 3.12、uv 以及 Kafka / JDBC / PostgreSQL 连接器 JAR。启动完成后你会得到 docker-compose.yml 中声明的四个服务服务暴露地址用途Redpandalocalhost:9092Kafka 兼容消息中间件承接 producer 写入与 Flink/consumer 读取Flink JobManagerhttp://localhost:8081作业协调者提供 Web UI 用于提交、监控与取消作业Flink TaskManager集群内部执行数据处理的 Worker默认 15 个任务槽、默认并行度 3PostgreSQLlocalhost:5432postgres/postgres结果落库供 SQL 查询验证如果你此前运行过 workshop 并残留了旧容器或数据卷务必做一次彻底的重置避免 topic 里的旧消息污染统计结果docker compose down -v docker compose build docker compose up -d注意容器名如workshop-redpanda-1假定目录名是workshop。若你重命名了目录后续所有docker exec命令中的容器名都要相应调整。数据准备Green Taxi 与字段筛选作业数据来自 NYC TLC 的绿色出租车行程文件2025 年 10 月读取 parquet 后只需保留以下 8 个字段lpep_pickup_datetime—— 上车时间datetime需转字符串lpep_dropoff_datetime—— 下车时间datetime需转字符串PULocationID—— 上车出租车区域 IDDOLocationID—— 下车出租车区域 IDpassenger_count—— 乘客数trip_distance—— 行程距离英里tip_amount—— 小费金额total_amount—— 总金额对比 workshop 使用的 producer.py只取 5 列、时间戳转 epoch 毫秒本作业字段更多、且要求把 datetime 列序列化为字符串以便 Flink 用TO_TIMESTAMP(..., yyyy-MM-dd HH:mm:ss)解析——这是本作业与 workshop 最关键的差异之一。Question 1确认 Redpanda 版本Redpanda 自带的 CLI 工具rpk用于管理 broker、topic 与消费组。在 Redpanda 容器内执行docker exec -it workshop-redpanda-1 rpk version输出会给出当前运行的 Redpanda 版本号。这道题的目的是确认环境就绪同时熟悉rpk这个后续创建/删除 topic 都要用到的工具。Question 2Producer——把数据写入 green-trips先创建作业专用 topicworkshop 中的ridestopic 由 broker 首次使用自动创建这里显式创建更稳妥docker exec -it workshop-redpanda-1 rpk topic create green-trips随后编写 producer读取 parquet → 筛选 8 列 → 每行转成 dict → 以 JSON 格式发送到green-trips。关键点是把两个 datetime 列先转成%Y-%m-%d %H:%M:%S格式字符串否则 JSON 序列化会失败或产生 Flink 无法解析的格式。整体骨架参照 workshop 的 producer.py改造如下import json import pandas as pd from time import time from kafka import KafkaProducer columns [ lpep_pickup_datetime, lpep_dropoff_datetime, PULocationID, DOLocationID, passenger_count, trip_distance, tip_amount, total_amount, ] df pd.read_parquet(green_tripdata_2025-10.parquet, columnscolumns) # datetime - 字符串与 Flink 的 yyyy-MM-dd HH:mm:ss 格式对齐 df[lpep_pickup_datetime] df[lpep_pickup_datetime].dt.strftime(%Y-%m-%d %H:%M:%S) df[lpep_dropoff_datetime] df[lpep_dropoff_datetime].dt.strftime(%Y-%m-%d %H:%M:%S) def json_serializer(data): return json.dumps(data).encode(utf-8) producer KafkaProducer( bootstrap_servers[localhost:9092], value_serializerjson_serializer, ) t0 time() for _, row in df.iterrows(): producer.send(green-trips, valuerow.to_dict()) producer.flush() t1 time() print(ftook {(t1 - t0):.2f} seconds)workshop 里发送 1000 行大约耗时 10 秒producer.py 中可见time.sleep(0.01)的节流逻辑本次是整月数据、全量发送用时量级取决于网络与机器。作业给出的候选答案是 10 / 60 / 120 / 300 秒用于判断你测得的时间落在哪个档位。Question 3Consumer——统计 trip_distance 5 的行程数写一个 Kafka 消费者读回全部消息。必须设置auto_offset_resetearliest否则新消费者默认从latest开始、只会看到订阅之后的新消息读不到已经写入的历史数据。group_id用于让 Kafka 记录该消费组的读取进度。参考 consumer.py 的写法import json from kafka import KafkaConsumer def json_deserializer(data): return json.loads(data.decode(utf-8)) consumer KafkaConsumer( green-trips, bootstrap_servers[localhost:9092], auto_offset_resetearliest, group_idgreen-trips-consumer, value_deserializerjson_deserializer, ) count 0 for message in consumer: if message.value[trip_distance] 5.0: count 1 print(ftrips with trip_distance 5: {count})作业的候选答案是 6506 / 7506 / 8506 / 9506取最接近你实际统计结果的选项。注意trip_distance在原始数据中为英里这里直接按数值比较不涉及单位换算。Part 2 预备PyFlink 作业的三个关键差异进入 PyFlink 部分前先明确与 workshop 代码aggregation_job.py的三处不同topic 名green-trips替代 workshop 的rides字段前缀datetime 列是lpep_前缀替代tpep_时间戳形态本作业的时间是字符串如2025-10-01 00:00:00而 workshop 中是 epoch 毫秒整数。因此源表 DDL 中要先用计算列把字符串转成 Flink 时间戳并声明 watermarklpep_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 SECONDTO_TIMESTAMP按给定格式解析字符串watermark 取最新事件时间 − 5 秒作为窗口发布结果的触发器容忍最多 5 秒的乱序到达workshop 的 README 对此有完整图解07-streaming/workshop/README.md。运行 Flink 作业前还需注意以下几点结果表先行在 PostgreSQL 中预先创建好各题的结果表用docker compose exec postgres psql -U postgres -d postgres或任意 SQL 客户端连接localhost:5432作业文件位置把你的.py作业文件放在workshop/src/job/下该目录被挂载进 Flink 容器容器内路径为/opt/src/job/提交命令docker exec -it workshop-jobmanager-1 flink run -py /opt/src/job/your_job.py并行度必须设为 1green-trips只有 1 个分区因此作业里要env.set_parallelism(1)。若并行度过高没有分配到数据的分区子任务idle consumer subtask不会消费消息会导致 watermark 无法推进、窗口迟迟不发布结果流式作业常驻运行Flink 流作业不会自行结束让作业跑 12 分钟、看到 PostgreSQL 里出现结果后可在 http://localhost:8081 的 Flink UI 上取消作业重复数据处理如果 producer 跑过多次topic 里会有重复消息。想彻底清空可删除并重建 topicdocker exec -it workshop-redpanda-1 rpk topic delete green-tripsQuestion 45 分钟 Tumbling Window 统计热门接客区需求统计每个PULocationID在每个 5 分钟滚动窗口内的行程数。先建结果表包含window_start、PULocationID、num_trips三列并建议像 workshop 的聚合表那样加上主键以启用 upsert 语义、让迟到数据能修正已发布的结果参见 aggregation_job.py 中PRIMARY KEY (...) NOT ENFORCED的用法CREATE TABLE trips_by_pulocation ( window_start TIMESTAMP, PULocationID INTEGER, num_trips BIGINT, PRIMARY KEY (window_start, PULocationID) );Flink 作业的核心 SQL 用TUMBLE函数第一个参数是源表第二个参数DESCRIPTOR(event_timestamp)必须指向定义了 watermark 的那一列第三个参数是窗口大小INSERT INTO trips_by_pulocation SELECT window_start, PULocationID, COUNT(*) AS num_trips FROM TABLE( TUMBLE(TABLE green_trips_source, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE) ) GROUP BY window_start, PULocationID;处理完所有数据后查询SELECT PULocationID, num_trips FROM trips_by_pulocation ORDER BY num_trips DESC LIMIT 3;取num_trips最大的PULocationID候选答案是 42 / 74 / 75 / 166。workshop 用 1 小时窗口做了同样的聚合演示只是多加了SUM(total_amount)与env.set_parallelism(3)见 aggregation_job.py你可以对照该文件理解 TUMBLE 的完整调用形态。Question 5Session Window 找最长会话需求以lpep_pickup_datetime作为事件时间、5 秒 watermark 容忍按PULocationID分组做会话窗口gap 为 5 分钟。会话窗口与 tumbling 的本质区别在于窗口大小不固定事件与上一事件间隔不超过 5 分钟就并入同一会话一旦出现超过 5 分钟的空档当前会话关闭。|--events--| gap(5min) |--events------| gap(5min) |--events--| | Session 1| | Session 2 | | Session 3|Flink SQL 中会话窗口写作SESSION(...)第一个参数仍是要按事件时间开窗的表INSERT INTO longest_sessions SELECT PULocationID, COUNT(*) AS num_trips FROM TABLE( SESSION(TABLE green_trips_source, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE) ) GROUP BY PULocationID, window_start, window_end;结果表自行定义至少包含PULocationID与num_trips然后找到单次会话内行程数最多的PULocationID即最长会话。候选答案是 12 / 31 / 51 / 81。workshop README 的Understanding window types一节对比了 tumbling / sliding / session 三种窗口的适用场景会话窗口常用于用户行为归因07-streaming/workshop/README.md这里把它用在出租车接客上语义同样成立。Question 61 小时 Tumbling Window 统计小费高峰需求用1 小时 tumbling 窗口聚合全量区域每小时的小费总额SUM(tip_amount)找出小费最高的那个小时。先建表CREATE TABLE hourly_tip_amounts ( window_start TIMESTAMP, total_tip_amount DOUBLE PRECISION );作业 SQL 与 Q4 同构只是窗口变为 1 小时、聚合函数换成SUM、且不再按PULocationID分组跨所有区域汇总INSERT INTO hourly_tip_amounts SELECT window_start, SUM(tip_amount) AS total_tip_amount FROM TABLE( TUMBLE(TABLE green_trips_source, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR) ) GROUP BY window_start;查询total_tip_amount最大的window_start候选答案是2025-10-01 18:00:00/2025-10-16 18:00:00/2025-10-22 08:00:00/2025-10-30 16:00:00。因为这里的结果表没有主键、且每小时只输出一行不会产生重复更新可以不用 upsert不过若担心同一窗口被重复发布仍可像 Q4 那样加主键。提交与检查清单完成全部 6 题后通过 cohort 提供的作业提交表单提交答案homework.yaml中声明了提交表单配置与截止时间见 cohorts/2027/07-streaming/homework.yaml。提交前建议对照以下清单自检Redpanda 版本号是否来自rpk version的真实输出producer 发送耗时与题面档位一致且flush()已调用未 flush 会导致部分消息滞留缓冲区、统计偏少consumer 统计用的是earliest偏移且只跑了一次重复运行同一group_id会从上一次位置继续导致计数不完整Flink 作业都设置了env.set_parallelism(1)watermark 能正常推进窗口结果已出现在 PostgreSQL 中再取消作业若 producer 重复发送过数据已用rpk topic delete green-trips清空后重新发送。与 workshop 源码的对照学习路径本作业的所有技术细节都能在仓库的 workshop 源码中找到对应实现建议按以下顺序对照研读07-streaming/workshop/docker-compose.yml四个服务的完整编排含 Redpanda 双监听地址、Flink 内存与并行度参数07-streaming/workshop/Dockerfile.flinkPyFlink 镜像构建过程与连接器 JAR 清单07-streaming/workshop/src/models.pyRide 数据模型与序列化/反序列化函数可据此改造成 green taxi 的 8 字段版本07-streaming/workshop/src/producers/producer.py 与 producer_realtime.py批量与实时两种生产方式后者还模拟了 ~20% 的 310 秒延迟事件用于观察 watermark 与 upsert 行为07-streaming/workshop/src/consumers/consumer.pyearliest偏移消费与反序列化07-streaming/workshop/src/job/pass_through_job.py最简 Kafka→JDBC 透传作业含latest-offset与TO_TIMESTAMP_LTZ的用法07-streaming/workshop/src/job/aggregation_job.pyTUMBLE 窗口 watermark upsert 主键的完整聚合作业是 Q4/Q6 的直接蓝本。把 workshop 中ridestopic、tpep_字段、epoch 毫秒时间戳这三处替换成本作业要求的green-tripstopic、lpep_字段与字符串时间戳再按上述要点调整并行度与结果表即可系统性地完成全部作业——这也正是以流处理思维迁移数据形态的实战训练目标。【免费下载链接】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),仅供参考
延伸阅读

更多相关文章

2026/9/12 7:20:03

企业级AI Agent竞争版图与技术落地:从MCP到LangGraph的实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 7:20:03

程序员如何避免氛围编程陷阱,提升工作效率

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 7:20:03

SpringBoot奶茶店管理系统开发实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 8:10:08

GEO产业图谱解析与企业竞争力分析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 8:10:08

RP2040低功耗本质:寄存器配置背后的物理级状态重构

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 8:10:08

SpringBoot部队战友管理系统设计与实现

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 8:10:08

磁单极近似与表面电荷密度方法在磁场计算中的对比

1. 磁单极近似与数值计算的基本概念在电磁学领域,计算磁场分布一直是个经典而富有挑战性的问题。当我第一次遇到长磁化圆柱体极尖间气隙的磁场计算需求时,发现教科书上的点磁单极近似方法虽然简单,但在实际工程应用中往往精度不足。这促使我深…

2026/9/12 8:10:08

基于MyEMS与LSTM的园区负荷预测实战:准确率95%

最近把公司园区能源管理平台上的负荷预测模块重新做了一版,底层用的是开源能源管理系统 MyEMS,预测模型用的是 LSTM 神经网络。最终在 2023 年全年留出的测试集上,MAPE 做到 4.7%,换算成大家常说的“准确率”大概在 95% 上下。要说…

2026/9/12 8:05:08

Biolaminin 521在PSC培养中的优化与应用

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 2:05:33

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/12 3:55:12

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/12 0:04:17

MATLAB仿生优化框架:长鼻浣熊算法多策略融合实现

简介:本资源是一份面向智能优化算法研究者与MATLAB初学者的仿生智能算法实践代码包,聚焦于长鼻浣熊优化算法(COA)的多策略改进与性能验证。针对传统COA易陷局部最优、收敛精度不足等问题,作者融合Circle映射初始化提升…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 JavaWeb 的校园一卡通管理系统的设计与实现 基于 JavaWeb 的校园卡业务管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 Java 的图书馆借阅管理平台的搭建与实现 基于 Java 的图书馆综合管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 6:29:36

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/10 15:19:50

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/12 6:37:43

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码