
Data Engineering Zoomcamp 07使用 Redpanda 与 kafka-python 构建生产级 Pub/Sub 流式处理基础【免费下载链接】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-zoomcampRedpanda 是一个与 Apache Kafka 协议兼容的流式数据平台本篇文章基于 Data Engineering Zoomcamp 第 07 模块的 Redpanda 实战示例从零搭建本地单节点集群并借助kafka-python实现出租车行程Taxi Ride数据的生产与消费。读者完成本指南后将掌握 cluster、broker、topic、producer、consumer group、offset 等流式核心概念并具备运行 rpk 命令行、访问 Redpanda Console UI 以及为 Capstone 流式项目打基础的全部技能。模块目标建立 Kafka/Redpanda 流式概念的地基在提交基于流式处理的 Capstone 项目之前需要先牢固掌握以下 Kafka/Redpanda 核心概念——本示例07-streaming/extras/python/redpanda_example/README.md正是围绕它们展开的clusters集群一组协同工作的 Redpanda broker 节点构成的整体是消息服务的基本运行单元brokers代理节点集群中实际存储数据、处理请求的节点本示例中的redpanda-1即为单节点集群中的一个 brokertopics主题消息按主题分类的逻辑通道生产者向主题写入、消费者从主题读取producers生产者将记录发布到指定主题的客户端consumers 与 consumer groups消费者与消费组从主题拉取记录并进行处理的客户端同一消费组内的消费者共同分担一个主题的分区负载data serialization and deserialization数据序列化与反序列化消息在网络上传输的是字节流客户端负责在发送前将对象序列化为字节、在接收后反序列化为对象replication and retention复制与保留主题分区在多个 broker 间的副本策略以及消息在集群中的保留时长/大小策略offsets偏移量记录在分区内的唯一序号消费者通过提交 offset 记录自己已消费到的位置。选择 Redpanda 作为入门载体有一个明显优势Redpanda 是一个单二进制镜像不像传统 Kafka 需要额外依赖 ZooKeeper因此用 Docker 一条命令即可拉起一个功能完整的流式集群非常适合用来学习 Kafka 协议体系。1. 前置准备安装 kafka-python如果你已经跟着模块 07 的视频完成了kafka-python的安装可以直接跳到 Docker 一节。否则这是本 Redpanda 课程唯一需要安装的 Python 包激活你的虚拟环境virtual environment执行安装命令pip install kafka-python仓库在 07-streaming/extras/python/requirements.txt 中锁定了本项目推荐使用的完整依赖版本kafka-python1.4.6 confluent_kafka requests avro faust fastavro其中本示例仅需kafka-python1.4.6其余依赖服务于同一模块下的 Avroconfluent-kafka、Faust 与 PySpark 流式示例。kafka-python提供KafkaProducer与KafkaConsumer两个核心类本示例中的 producer.py 与 consumer.py 均直接基于它们构建。2. Docker启动单节点 Redpanda 集群进入示例目录并启动集群cd 07-streaming/extras/python/redpanda_example/ docker-compose up -d对应的 docker-compose.yaml 定义了三个部分1Redpanda 主节点redpanda-1镜像版本docker.redpanda.com/redpandadata/redpanda:v23.2.26。启动命令中的关键参数参数取值作用--smp 11限制为单 CPU 核适合本地学习环境节省资源--reserve-memory 0M0M不为内部用途预留额外内存降低本地运行门槛--overprovisioned开关允许在资源不足的机器上运行过度供给模式学习场景常用--node-id 11指定当前节点在集群中的 ID--kafka-addrPLAINTEXT://0.0.0.0:29092,OUTSIDE://0.0.0.0:9092分别绑定容器内PLAINTEXT与外部OUTSIDE的 Kafka API 监听地址--advertise-kafka-addrPLAINTEXT://redpanda-1:29092,OUTSIDE://localhost:9092对外广播的地址容器内部客户端通过redpanda-1:29092访问宿主机客户端通过localhost:9092访问--pandaproxy-addr/--advertise-pandaproxy-addr0.0.0.0:8082/0.0.0.0:28082等HTTP ProxyPandaProxy与 Admin API 的绑定与广播地址--rpc-addr/--advertise-rpc-addr0.0.0.0:33145/redpanda-1:33145节点间内部 RPC 通信地址端口映射中值得注意9092:9092供宿主机上的 Python 客户端连接对应BOOTSTRAP_SERVERS [localhost:9092]9644:9644是 Redpanda 的 Admin API 端口Console UI 依赖它展示集群信息8082:8082是 HTTP Proxy 端口。2被注释掉的第二节点redpanda-2如果你想体验两节点集群取消这段注释即可对应--node-id 2与--seeds redpanda-1:33145它演示了通过--seeds参数让新节点加入既有集群的组网方式。注意其镜像版本为v23.1.1若要与主节点保持一致可自行调整。3Redpanda Consoleredpanda-console镜像console:v2.2.2一个 Web 管理界面通过环境变量注入配置——kafka.brokers: [redpanda-1:29092]指向容器内 Kafka 地址redpanda.adminApi.urls: [http://redpanda-1:9644]指向 Admin API端口8080:8080暴露给浏览器访问。提示容器内部署后Python 示例使用宿主机地址localhost:9092OUTSIDE 监听而 Console 使用容器内地址redpanda-1:29092PLAINTEXT 监听这正是 docker-compose 中配置两套 advertise 地址的原因。3. 配置 rpk 命令别名Redpanda 自带的命令行工具rpkRedpanda keeper已经内置在 Docker 镜像中。为了不必每次打开容器交互终端在宿主机终端设置如下别名alias rpkdocker exec -ti redpanda-1 rpk rpk version执行后可以看到类似v23.2.26 (rev 328d83a06e)的版本号。Redpanda 遵循major.minor[.build[.revision]]的语义化版本规则其中主要版本号v23是关键——保证你得到的结果与本仓库文档演示一致。[!TIP] 如果你在 2024 年 3 月之后阅读本文且希望将 Docker 配置升级到更新的 Redpanda 镜像可访问 Docker Hub 上的vectorized/redpandatags 页面将新版本号替换到 docker-compose.yaml 的镜像标签中即可。4. 运行 Kafka 生产者 - 消费者示例在并排的两个终端标签页中打开两个 shell并在每个终端中激活虚拟环境然后分别执行# 在第 1 个终端标签页启动消费者 python -m consumer.py # 在第 2 个终端标签页启动生产者 python -m producer.py多次重复执行python -m producer.py可以观察到消费者终端会实时、自动地消费到新产生的events消息记录。4.1 生产者实现剖析producer.pyproducer.py 中的JsonProducer继承自KafkaProducer工作流程分三步读取 CSV 记录read_records()打开 rides.csv18 列、265 行数据的纽约出租车样例跳过表头后将每行封装为一个Ride对象配置序列化器key_serializer将分区键pu_location_id上车地点 ID转为字节value_serializer将Ride对象的__dict__经json.dumps(..., defaultstr)转为 JSON 字节串——Kafka 协议要求键值对以二进制形式在网络中传输发布消息publish_rides()通过producer.send(topic..., key..., value...)异步发送并通过record.get()等待写入结果输出每条记录落盘的 offsetconfig { bootstrap_servers: BOOTSTRAP_SERVERS, key_serializer: lambda key: str(key).encode(), value_serializer: lambda x: json.dumps(x.__dict__, defaultstr).encode(utf-8) }4.2 消费者实现剖析consumer.pyconsumer.py 中的JsonConsumer封装了KafkaConsumer核心配置包括config { bootstrap_servers: BOOTSTRAP_SERVERS, auto_offset_reset: earliest, # 无已提交 offset 时从最早消息开始消费 enable_auto_commit: True, # 自动提交消费进度 key_deserializer: lambda key: int(key.decode(utf-8)), value_deserializer: lambda x: loads(x.decode(utf-8), object_hooklambda d: Ride.from_dict(d)), group_id: consumer.group.id.json-example.1, # 消费组标识 }consume_from_kafka()中consumer.subscribe(topics)订阅主题后进入while True轮询循环consumer.poll(1.0)以 1 秒为超时拉取一批消息注释说明 SIGINT 信号无法在 poll 阻塞期间被处理因此限制超时时间以保证CtrlC能及时中断拉取到的消息按「键 值」打印。key_deserializer将字节键还原为整数value_deserializer借助 ride.py 的Ride.from_dict()将 JSON 恢复为Ride对象。4.3 Ride 领域模型与数据形态ride.py 定义了Ride类其构造函数接收 18 个字段的列表覆盖样例数据 rides.csv 的全部列vendor_id、tpep_pickup_datetime、tpep_dropoff_datetime按%Y-%m-%d %H:%M:%S解析、passenger_count、trip_distance、rate_code_id、store_and_fwd_flag、pu_location_id、do_location_id、payment_type以及fare_amount、extra、mta_tax、tip_amount、tolls_amount、improvement_surcharge、total_amount、congestion_surcharge等金额字段使用Decimal精确保存。from_dict()类方法提供了反序列化入口与消费者端的object_hook配合形成闭环。关键连接配置集中在 settings.pyINPUT_DATA_PATH ../resources/rides.csv BOOTSTRAP_SERVERS [localhost:9092] KAFKA_TOPIC rides_jsonINPUT_DATA_PATH指向样例数据文件相对本文件所在目录BOOTSTRAP_SERVERS即 Kafka API 的引导地址与 docker-compose 中OUTSIDE://localhost:9092的广播地址一致KAFKA_TOPIC为rides_json生产与消费双方必须使用相同的主题名才能互通。4.4 JSON 无 Schema 的隐患consumer.py 文件末尾的注释点明了 JSON 格式的软肋JSON 本身不携带 Schema因此当字段被删除、新增或数据类型变化时Ride类与消息传输依然可以“顺畅”运行但下游分析端的数据集会缺失该列导致仪表盘失败进而侵蚀对数据与流程的信任。这正是同一模块下引入 Avro 序列化见 avro_example的动机——用强 Schema 约束来保护下游消费契约。5. 通过 Redpanda Console UI 观察集群浏览器访问 http://localhost:8080即可打开 Redpanda Console Web 界面直观查看集群列表与健康状态各 broker 节点信息Topicsrides_json主题的分区partitions、副本replicas、消息数、保留策略等元数据Consumer Groupsconsumer.group.id.json-example.1消费组的成员与 lag消费滞后情况。Console 的配置在 docker-compose.yaml 的redpanda-console服务中定义其中schemaRegistry.enabled: false表示本示例未启用 Schema Registry仅使用 JSON。6. rpk 命令速查表rpk是管理 Redpanda 集群的核心命令行工具常用命令如下完整用法可参考 Redpanda 官方 get-started-rpk 博客文章# 设置 rpk 别名已在第 3 节配置 alias rpkdocker exec -ti redpanda-1 rpk # 获取集群信息broker 列表、版本、分区情况 rpk cluster info # 创建主题 topic_name指定 m 个分区、n 个副本因子 rpk topic create [topic_name] --partitions m --replicas n # 列出主题普通列表 / 详细列表含分区、副本等细节 rpk topic list rpk topic list --detailed # 查看指定主题的配置 rpk topic describe [topic_name] # 从指定主题消费消息实时打印到终端 rpk topic consume [topic_name] # 列出集群中的消费组 rpk group list # 查看某个消费组的详细信息例如上一步查到的 my-group rpk group describe my-group这几个命令覆盖了主题生命周期管理与消费组观测两条主线create / list / describe管理主题consume直接验证生产链路group list / describe则用于排查消费组 lag 或成员分配问题。配合rpk cluster info可以对标 第 4 节 中 Python 生产/消费的结果做交叉验证。7. 更进一步的流式学习资源如果你希望在模块 07 基础上继续夯实流式基础可关注 Redpanda University需注册 Redpanda 账号课程免费RP101: Getting Started with Redpanda——Redpanda 上手实践RP102: Stream Processing with Redpanda——基于 Redpanda 的流式处理SF101: Streaming Fundamentals——流式处理基础理论SF102: Kafka building blocks——Kafka 核心构件。如果你自认为已具备扎实的流式与 Kafka 基础可以跳过这些补充课程。此外本仓库还提供了同一主题的进阶扩展PySpark 与 Redpanda 结合的流式处理见 streams-example/redpanda/README.md以及完整的 PyFlink 实时流式管道 workshop见 07-streaming/workshop。总结本文以 Data Engineering Zoomcamp 第 07 模块的 Redpanda 示例为骨架完成了从集群搭建、rpk 工具链、生产者/消费者代码到 Console UI 观测的完整闭环docker-compose up -d一键拉起集群python -m producer.py将 rides.csv 中的出租车数据发布到rides_json主题python -m consumer.py实时拉取并反序列化。通过这套最小可运行示例你将集群、broker、topic、producer、consumer group、序列化、复制保留与 offset 这些 Kafka/Redpanda 概念逐一落到可操作的实践中为后续的流式 Capstone 项目奠定坚实的地基。【免费下载链接】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),仅供参考