ARTICLE DETAIL

资讯详情

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

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

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),仅供参考
返回列表