Data Engineering Zoomcamp 流式处理:用 Python 消费 Kafka 消息——从反序列化到 Consumer Group 实战
2026/9/12 4:24:24 网站建设 项目流程

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_idauto_offset_resetvalue_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 是 Redpanda(redpandadata/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_deserializer

Kafka 消费者拿到的是原始字节。最朴素的做法是:先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.dumpsencode('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(PULocationID=186, DOLocationID=79, trip_distance=1.72, # total_amount=17.31, tpep_pickup_datetime=1730429702000)

验证通过后,ride_deserializer可以直接作为value_deserializer传给KafkaConsumer——Kafka 客户端会在每条消息到达时自动调用它,于是message.value直接就是Ride对象,消费代码里不再需要任何手工转换。

连接 Kafka:KafkaConsumer核心参数

现在创建消费者连接。仓库中的完整实现位于 consumer.py:

from kafka import KafkaConsumer server = 'localhost:9092' topic_name = 'rides' consumer = KafkaConsumer( topic_name, bootstrap_servers=[server], auto_offset_reset='earliest', group_id='rides-console', value_deserializer=ride_deserializer )

逐参数拆解:

  • bootstrap_servers:broker 接受连接的地址。localhost:9092是因为我们在宿主机(Docker 外部)运行。若多个 broker 可传列表,如['kafka1:9092', 'kafka2:9092']——客户端通过 bootstrap 获取集群元数据后,会连向 broker 返回的 advertised 地址进行实际数据传输(Redpanda 的双监听地址设计见 02-redpanda.md)。
  • auto_offset_reset='earliest':决定新消费组(该 topic 无已提交 offset)从何处开始读:
    • 'earliest':从 topic 开头重放所有历史消息;
    • 'latest'(kafka-python 默认):只消费连接建立之后到达的新消息。
  • group_id='rides-console':标识消费组。Kafka 按 (group, partition) 维度记录每个组已消费到的 offset,因此用同一 group_id 重启消费者会从上次的位置继续,而不是重复消费;换一个新 group_id 则相当于"新人",从头(按auto_offset_reset规则)读起。
  • value_deserializer:每条消息 value 的字节 → 对象转换函数,即上一节定义的ride_deserializer

依赖提醒:kafka-python由项目 pyproject.toml 声明(kafka-python>=2.3.0),与pandaspyarrowpsycopg2-binary一并由 uv 管理。

消费循环:把事件打印出来

KafkaConsumer是可迭代对象,for message in consumer会阻塞等待新消息。由于value_deserializer已把 value 变成Ride,循环体可以专注于业务处理:

from datetime import datetime print(f"Listening to {topic_name}...") count = 0 for message in consumer: ride = message.value pickup_dt = datetime.fromtimestamp(ride.tpep_pickup_datetime / 1000) print(f"Received: PU={ride.PULocationID}, DO={ride.DOLocationID}, " f"distance={ride.trip_distance}, amount=${ride.total_amount:.2f}, " f"pickup={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 >= 10break,避免控制台无限刷屏——这是调试流式消费者的常用手法。
  • 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=$..., pickup=2025-... ... ... received 10 messages so far (stopping after 10 for demo)

由于设置了auto_offset_reset='earliest'且是首次消费,消费者会从头开始重放ridestopic 中已有的消息(例如生产者刚写入的 1000 条出租车行程),打印 10 条后退出。

源码级对照:consumer.py的完整调用链

仓库中的 consumer.py 与上面逐段讲解的代码一一对应,只有两处工程化补充:

  1. sys.path.insert(0, str(Path(__file__).parent.parent)):把src/加入模块搜索路径,使from models import ride_deserializer直接可用。这个约定贯穿生产者、消费者与后续 Flink job 的所有脚本(见 producer.py 与 consumer_postgres.py)。
  2. 复用models模块:反序列化逻辑不散落在各处,而是集中在 models.py 的ride_deserializer,任何消费者(控制台、PostgreSQL、Flink)都用同一份解析逻辑,保证 schema 一致性。

从源码结构看,该目录下的消费端脚本呈渐进式设计:consumer.py(打印)→consumer_postgres.py(落库)→job/下的 Flink 作业(窗口聚合),同一套反序列化与消费模型被逐级复用。

进阶语义:Consumer Group 与 offset 的行为差异

理解了group_id之后,一个关键问题自然浮现:earliestlatest到底在什么时机生效?答案是仅在消费组对该 topic 无已提交 offset(或 commit 无效)时

  • 首次以rides-console消费 +earliest→ 重放全部历史消息;
  • 同组再次启动 → 从上次提交的 offset 继续,auto_offset_reset不再生效;
  • 换新组名(如rides-to-postgres)→ 视为全新组,再次从earliestlatest起步。

这一语义在后续 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的消费骨架完全一致,仅将打印替换为参数化 INSERT:

cur.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.yml),Docker 内部服务应改用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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询