Data Engineering Zoomcamp 流式处理:用 Python 消费 Kafka 消息——从反序列化到 Consumer Group 实战

发布时间:2026/9/12 4:24:46

Data Engineering Zoomcamp 流式处理:用 Python 消费 Kafka 消息——从反序列化到 Consumer Group 实战 Data Engineering Zoomcamp 流式处理用 Python 消费 Kafka 消息——从反序列化到 Consumer Group 实战【免费下载链接】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 流式处理模块PyFlink Stream Processing Workshop中的一个关键环节如何用 Python 从 Kafka本项目实际以 Redpanda 作为兼容实现消费 NYC 黄色出租车事件流。你将掌握 Kafka 字节消息的反序列化思路、KafkaConsumer核心参数group_id、auto_offset_reset、value_deserializer的语义并基于仓库中的完整源码跑通一个可持续扩展的消费端脚本为后续 Flink 流处理与 PostgreSQL 落库打好基础。背景Kafka 消费模型——字节进字节出在 PyFlink: Stream Processing Workshop 中整条实时管道被构建为Producer (Python) - Kafka (Redpanda) - Flink - PostgreSQL消费端是这条管道承上启下的枢纽生产者在 03-produce-messages-to-kafka.md 中把 DataFrame 行序列化成 JSON 字节并写入ridestopic而消费端则要完成对称的逆操作——把 Kafka 交付的原始字节还原成结构化对象。Kafka 协议本身对消息内容零假设它只负责存储和分发字节数组byte array。所有语义是 JSON、Avro 还是 Protobuf都由客户端自行编解码。因此写一个消费端的第一件事不是连 broker而是先想清楚字节如何变成我代码里的对象。说明本 workshop 中所有 Kafka 均指 Kafka 协议与概念底层 broker 是 Redpandaredpandadata/redpanda:v25.3.9配置细节见 02-redpanda.md 与 docker-compose.yml。任何 Kafka 客户端库无需任何改动即可对接。共享数据模型Ridedataclass消费者收到的每个事件对应一次出租车行程。为了让消息具备明确 schema项目在 models.py 中定义了Ride数据类from dataclasses import dataclass dataclass class Ride: PULocationID: int DOLocationID: int trip_distance: float total_amount: float tpep_pickup_datetime: int # epoch milliseconds要点解析tpep_pickup_datetime是整数epoch 毫秒而非字符串——这是与 Flink 协作的关键约定。生产者侧通过int(row[tpep_pickup_datetime].timestamp() * 1000)把 pandas Timestamp 转成毫秒时间戳见 producer.py消费端取到毫秒数后由业务代码决定何时转成可读时间。生产端与消费端共用同一份models.py。该文件同时定义了序列化ride_from_row与反序列化ride_deserializer两侧的工具函数这正是schema boundary 显式化的体现表格式输入 → 事件 → 字节 → topic 记录 → 字节 → 对象。一步到位的反序列化ride_deserializerKafka 消费者拿到的是原始字节。最朴素的做法是先decode(utf-8)成 JSON 字符串 →json.loads成 dict → 再手动构造Ride(**ride_dict)。每次都写这三步很繁琐因此本项目把它封装成一个函数一步完成解码 解析 构造对象import json def ride_deserializer(data): json_str data.decode(utf-8) ride_dict json.loads(json_str) return Ride(**ride_dict)这正好是生产者侧ride_serializerdataclasses.asdict(ride)→json.dumps→encode(utf-8)的镜像操作环节生产者消费者对象 ↔ 字典dataclasses.asdict(ride)Ride(**ride_dict)字典 ↔ 字符串json.dumps(ride_dict)json.loads(json_str)字符串 ↔ 字节json_str.encode(utf-8)data.decode(utf-8)用样例字节验证反序列化Kafka 交付给你的就是编码后的二进制字符串。可以用一段样例 JSON 字节来验证函数行为这正是 Kafka 中消息的真实形态test_bytes json.dumps({ PULocationID: 186, DOLocationID: 79, trip_distance: 1.72, total_amount: 17.31, tpep_pickup_datetime: 1730429702000 }).encode(utf-8) ride_deserializer(test_bytes) # Ride(PULocationID186, DOLocationID79, trip_distance1.72, # total_amount17.31, tpep_pickup_datetime1730429702000)验证通过后ride_deserializer可以直接作为value_deserializer传给KafkaConsumer——Kafka 客户端会在每条消息到达时自动调用它于是message.value直接就是Ride对象消费代码里不再需要任何手工转换。连接 KafkaKafkaConsumer核心参数现在创建消费者连接。仓库中的完整实现位于 consumer.pyfrom kafka import KafkaConsumer server localhost:9092 topic_name rides consumer KafkaConsumer( topic_name, bootstrap_servers[server], auto_offset_resetearliest, group_idrides-console, value_deserializerride_deserializer )逐参数拆解bootstrap_serversbroker 接受连接的地址。localhost:9092是因为我们在宿主机Docker 外部运行。若多个 broker 可传列表如[kafka1:9092, kafka2:9092]——客户端通过 bootstrap 获取集群元数据后会连向 broker 返回的 advertised 地址进行实际数据传输Redpanda 的双监听地址设计见 02-redpanda.md。auto_offset_resetearliest决定新消费组该 topic 无已提交 offset从何处开始读earliest从 topic 开头重放所有历史消息latestkafka-python 默认只消费连接建立之后到达的新消息。group_idrides-console标识消费组。Kafka 按 (group, partition) 维度记录每个组已消费到的 offset因此用同一 group_id 重启消费者会从上次的位置继续而不是重复消费换一个新 group_id 则相当于新人从头按auto_offset_reset规则读起。value_deserializer每条消息 value 的字节 → 对象转换函数即上一节定义的ride_deserializer。依赖提醒kafka-python由项目 pyproject.toml 声明kafka-python2.3.0与pandas、pyarrow、psycopg2-binary一并由 uv 管理。消费循环把事件打印出来KafkaConsumer是可迭代对象for message in consumer会阻塞等待新消息。由于value_deserializer已把 value 变成Ride循环体可以专注于业务处理from datetime import datetime print(fListening to {topic_name}...) count 0 for message in consumer: ride message.value pickup_dt datetime.fromtimestamp(ride.tpep_pickup_datetime / 1000) print(fReceived: PU{ride.PULocationID}, DO{ride.DOLocationID}, fdistance{ride.trip_distance}, amount${ride.total_amount:.2f}, fpickup{pickup_dt}) count 1 if count 10: print(f\n... received {count} messages so far (stopping after 10 for demo)) break consumer.close()值得注意的实现细节毫秒时间戳的换算ride.tpep_pickup_datetime是 epoch 毫秒10^13 量级除以 1000 得到秒datetime.fromtimestamp才能正确解释若直接传毫秒会导致年份错误。演示限流count 10时break避免控制台无限刷屏——这是调试流式消费者的常用手法。consumer.close()显式关闭消费者释放网络连接与本地 offset 状态。完整的循环写法中应放在finally里或使用上下文管理器保证异常时也能正确关闭。运行消费者uv run python src/consumers/consumer.py若从仓库起步目录为cohorts/2027/07-streaming/code对应文件是 consumer.py。运行前请确保已按 02-redpanda.md 启动 Redpanda并按 03-produce-messages-to-kafka.md 先运行生产者向ridestopic 写入数据。预期输出Listening to rides... Received: PU..., DO..., distance..., amount$..., pickup2025-... ... ... received 10 messages so far (stopping after 10 for demo)由于设置了auto_offset_resetearliest且是首次消费消费者会从头开始重放ridestopic 中已有的消息例如生产者刚写入的 1000 条出租车行程打印 10 条后退出。源码级对照consumer.py的完整调用链仓库中的 consumer.py 与上面逐段讲解的代码一一对应只有两处工程化补充sys.path.insert(0, str(Path(__file__).parent.parent))把src/加入模块搜索路径使from models import ride_deserializer直接可用。这个约定贯穿生产者、消费者与后续 Flink job 的所有脚本见 producer.py 与 consumer_postgres.py。复用models模块反序列化逻辑不散落在各处而是集中在 models.py 的ride_deserializer任何消费者控制台、PostgreSQL、Flink都用同一份解析逻辑保证 schema 一致性。从源码结构看该目录下的消费端脚本呈渐进式设计consumer.py打印→consumer_postgres.py落库→job/下的 Flink 作业窗口聚合同一套反序列化与消费模型被逐级复用。进阶语义Consumer Group 与 offset 的行为差异理解了group_id之后一个关键问题自然浮现earliest与latest到底在什么时机生效答案是仅在消费组对该 topic 无已提交 offset或 commit 无效时首次以rides-console消费 earliest→ 重放全部历史消息同组再次启动 → 从上次提交的 offset 继续auto_offset_reset不再生效换新组名如rides-to-postgres→ 视为全新组再次从earliest或latest起步。这一语义在后续 09-offsets-earliest-vs-latest.md 中被正式化为 Flink 的scan.startup.mode三档latest-offset只读新消息生产常用、earliest-offset重放历史用于回填/重算、timestamp从指定时刻恢复故障恢复场景。两处表述一致可见本模块把消费起点作为一条贯穿 Python 与 Flink 的主线知识。多消费组的实际意义在 05-save-events-to-postgresql.md 中项目刻意让控制台消费者与 PostgreSQL 消费者使用不同group_idrides-consolevsrides-to-postgres这样两者各自独立跟踪 offset都能读到全部消息——这正是 Kafka 消费组模型的核心价值不同下游互不干扰各自维护进度。延伸从打印到落库打印只是调试手段。把消费者升级为数据持久化只需两步在 docker-compose.yml 中加入 PostgreSQL 服务并新建 consumer_postgres.py。该脚本与consumer.py的消费骨架完全一致仅将打印替换为参数化 INSERTcur.execute( INSERT INTO processed_events (PULocationID, DOLocationID, trip_distance, total_amount, pickup_datetime) VALUES (%s, %s, %s, %s, %s), (ride.PULocationID, ride.DOLocationID, ride.trip_distance, ride.total_amount, pickup_dt) )这也直接点出了手写消费者的边界窗口聚合、崩溃恢复、并行分区分配、多 sink 支持都需要自行实现——这正是后续引入 Flink 的动机详见 05-save-events-to-postgresql.md 末尾的讨论。常见问题排查现象可能原因与对策消费者启动后收不到任何消息broker 未启动docker compose up redpanda -d或生产者还没运行、topic 为空或auto_offset_reset设为latest而消息在连接前已写入重启后从头重复消费用了新的group_id或上一进程未提交 offset 即被终止时间显示年份异常datetime.fromtimestamp()收到的仍是毫秒值需先/ 1000连接失败localhost:9092确认端口映射9092:9092存在见 docker-compose.ymlDocker 内部服务应改用redpanda:29092小结消费 Kafka 消息的完整套路可浓缩为四步定义数据模型Ride→ 编写字节到对象的反序列化函数ride_deserializer→ 配置消费者KafkaConsumer四要素→ 循环处理message.value。在 Data Engineering Zoomcamp 的这条实时管道中这个消费端既是验证生产数据的探针也是通往 PostgreSQL 落库与 Flink 窗口计算的起点。完整的可直接运行代码见 code/src/consumers/consumer.py其后续演进路线落库、Flink、offset 语义可在 05-save-events-to-postgresql.md 与 09-offsets-earliest-vs-latest.md 中继续研读。【免费下载链接】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 4:24:46

deepin与Windows 11双系统安装实战:从磁盘规划到引导修复

说实话,双系统这碗饭我吃了很多年,Windows和Linux来回切换的需求一直没断过。最近腾出一台主力笔记本做测试,直接把 deepin 装上,和 Windows 11 组了双系统。整个流程走下来,比早年装 Ubuntu 双系统顺滑不少&#xff0…

2026/9/12 4:24:46

Node.js+Vue全栈在线考试系统开发实践

1. 项目概述:基于Node.jsVue的全栈在线考试系统这个项目本质上是一个融合了前后端技术的教育信息化解决方案。我在实际开发中发现,传统的纸质考试或单机版考试系统存在诸多痛点:试卷批改效率低、成绩统计滞后、作弊风险高、数据分析困难。而采…

2026/9/12 4:44:49

MindSpore多模态大模型产线落地实战:昇腾边缘实时推理优化

1. 项目概述:这不是又一个“跑通Demo”的故事,而是把多模态大模型真正焊进产线的实操笔记我做AI工程落地快八年了,从最早用TensorFlow 1.x搭CV pipeline,到后来在华为昇腾集群上跑通第一个千亿参数大模型推理服务,踩过…

2026/9/12 4:44:49

Linux自定义Shell开发指南:从基础架构到高级功能实现

/* 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 4:44:49

视觉模型边缘部署实战:延迟预算、量化与断网兜底

把视觉模型从云端挪到现场边缘盒子这件事,我前后做过三套,最早的版本踩的坑最多:摄像头往云上推流,云端跑检测,结果返回结果那一下总是慢半拍,机械臂抓偏、AGV 刹不住、直播里的识别框永远追不上人。这套路…

2026/9/12 4:44:49

STM32开发三大深坑:工程配置、时钟系统与外设调试

/* 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 4:44:48

LLM API Gateway生产落地:自建方案与关键增强点全解析

搞LLM应用最头疼的事,不是模型效果不够好,而是怎么让服务在生产环境里真正稳定跑起来。你开发的时候用Python脚本直连OpenAI或者其他模型API挺爽,参数随手一调,请求一发,结果就回来了。但一旦要上生产,面对…

2026/9/12 4:39:48

Python tkinter Text组件选择事件深度解析与应用

1. 深入理解tkinter的Text组件与虚拟事件机制在Python GUI开发领域&#xff0c;tkinter作为标准库中的"常青树"&#xff0c;其Text组件堪称构建文本编辑功能的瑞士军刀。而<<Selection>>这个看似简单的虚拟事件&#xff0c;实则是处理文本选择操作的关键…

2026/9/12 2:05:33

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

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

2026/9/12 3:55:12

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

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

2026/9/9 16:31:09

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

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

2026/9/12 0:04:17

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

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

2026/9/12 0:04:17

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

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

2026/9/12 0:04:17

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

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

2026/9/10 12:32:02

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

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

2026/9/10 15:19:50

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

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

2026/9/10 15:49:53

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

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

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

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

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