
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),仅供参考