消息中间件:Kafka 快速入门
kafkahttps://github.com/apache/kafkakafka-pythonhttps://github.com/dpkp/kafka-python1、Kafka 简介什么是 kafkaApache Kafka 是一款分布式、高吞吐、可持久化、分区、多副本的消息队列消息中间件基于发布 / 订阅模式。 典型场景日志收集、实时数据流系统解耦、异步通信流量削峰填谷实时计算Flink/Spark Streaming 数据源Kafka 为什么这么强 高吞吐磁盘顺序读写 批量压缩。“Kafka 基于磁盘存储为什么能比内存队列更快” 答案在于磁盘顺序读写和批量压缩。传统磁盘 IO 慢的核心是 “随机读写”磁头频繁移动而 Kafka 中每个 Partition 的数据是按 “日志文件” 形式顺序存储的 ——Producer 写入时只能在文件末尾追加Append OnlyConsumer 读取时也按顺序从头往后读避免了随机 IO 的开销磁盘顺序读写的速度甚至接近内存。同时Kafka 支持对消息进行批量压缩如 Gzip、SnappyProducer 将多个小消息打包压缩后发送减少网络传输量Broker 存储压缩后的消息减少磁盘占用Consumer 读取后解压处理整体提升了端到端的吞吐效率。高可靠副本机制 数据持久化。数据丢失是分布式系统的 “噩梦”而 Kafka 通过副本机制和数据持久化确保数据不丢失。副本机制每个 Partition 会有多个副本副本数可配置通常为 3包括 1 个 Leader 和多个 Follower。Leader 负责读写Follower 实时同步 Leader 的数据。当 Leader 宕机时Kafka 会从 Follower 中选举新的 Leader确保数据不中断同时通过 “min.insync.replicas” 参数可配置 “最少同步副本数”如设置为 2需 Leader 和至少 1 个 Follower 同步成功才算消息写入成功进一步降低数据丢失风险。数据持久化所有消息会被持久化到磁盘即使 Broker 重启数据也不会丢失。同时Kafka 支持配置消息的 “保留时间”如 7 天过期数据会自动删除避免磁盘占满。低延迟零拷贝技术。在数据从 Broker 发送到 Consumer 的过程中传统方式需要经过 “磁盘 → 内核缓冲区 → 用户缓冲区 → 内核 Socket 缓冲区 → 网卡” 多个步骤存在多次数据拷贝耗时较长。而 Kafka 采用零拷贝技术通过 Linux 的 sendfile 系统调用直接将磁盘文件数据映射到内核缓冲区再从内核缓冲区发送到网卡跳过了用户态的拷贝步骤将数据传输延迟降低到毫秒级满足实时计算场景的需求。Kafka 架构在整个 Kafka 集群中 Producer 将消息发送给 broker然后 broker 再将接收到的消息存储到磁盘中然后 Consumer 再从 Broker 订阅并消费消息。ZooKeeper 则是 Kafka 集群用来负责集群元数据的管理、控制器的选举等操作的。Kafka 集群结构图Kafka Raft 原理 结构图Kafka 物理存储结构图Kafka 性能优化一个典型的 Kafka 架构会包括 Producer、broker、Cosumer 等角色以及一个 ZooKeeper 集群核心组件Topic主题消息的逻辑分类可以理解为一个队列。生产者Producer发送消息时必须指定 Topic消费者Consumer订阅消息也需基于 Topic。生产者和消费者面向的都是一个 topic。Partition分区Topic 可物理拆分为多个分区Partition一个分区对应磁盘上一个文件夹也叫主题分区分区分散部署在集群不同 Broker是 Kafka 实现高吞吐、横向扩展的核心。 每个分区内部消息有序但 Topic 全局无序。生产者依据分区策略路由消息携带 key 时通过hash(key) % 分区数分配相同 key 消息固定进入同一分区无 key 则采用轮询分发。多分区架构支持生产者并行写入、消费者并行拉取有效提升整体吞吐量。同一分区的不同副本中保存的信息是相同的通过多副本机制实现了故障的自动转移当集群中某个 broker 失效时仍然能保证服务可用可以提升容灾能力。Broker是 Kafka 服务节点每个 Broker 就是一台运行 Kafka 服务的机器或者独立的进程一个 Kafka 集群由多个 Broker 组成一个 broker 可以容纳多个 topic。集群中 Broker 的数量决定了系统的容灾能力通常建议至少 3 个 Broker 组成集群避免单点故障。Broker 有 “首领”Leader和 “追随者”Follower之分Leader 负责处理 Topic 的读写请求Follower 仅同步 Leader 的数据当 Leader 故障时Follower 会通过选举机制成为新的 Leader保证服务连续性。Replica副本为保证集群中某个节点发生故障时该节点上的 partition 数据不丢失且 Kafka 仍然可以继续工作Kafka 提供了副本机制。一个分区会有多个副本副本之间是一主( Leader读写)多从(Follower只同步数据)的关系Leader 对外提供服务而 Follower 只是被动地同步 Leader 不对外提供服务。Offset偏移量分区内消息唯一序号消费者依靠 offset 记录消费位置Producer生产者消息生产者就是发送消息到 Topic 的客户端。发送时Producer 会根据一定的策略如轮询、按消息 key 哈希将消息分配到 Topic 的不同 Partition确保数据在 Partition 间均匀分布。同时Producer 支持 “acks” 参数配置0不等待确认1等待 Leader 确认-1等待 Leader 和所有 Follower 确认平衡数据可靠性与发送效率。Consumer消费者消息消费者从 Topic 拉取消息的客户端。Consumer 加入消费组从指定分区拉取消息处理完成后提交 offset。与 Producer 不同的是 Consumer 必须属于一个Consumer Group消费者组—— 同一 Consumer Group 中的多个 Consumer 会分工读取 Topic 的不同 Partition一个 Partition 只能被同一 Group 中的一个 Consumer 消费避免重复消费而不同 Consumer Group 可独立消费同一 Topic 的数据实现 “一份数据多端处理” 的场景如一份订单数据既用于实时计算也用于离线存储。Consumer Group消费组多个消费者共同组成一个组目的是让多个消费者同时消费同一个 Topic 中的消息可以加速整个消费者端的吞吐量。同一个分区只能被组内一个消费者消费。消费者组间互不影响。所有的消费者都属于某个消费者组即消费者组是逻辑上的一个订阅者。消费规则消费者数量 ≤ 分区数量消费者多于分区多余消费者空闲leader每个分区多个副本的 ” 主 “生产者发送数据的对象以及消费者消费数据时的对象都是 leader。follower每个分区多个副本的 “从”实时从 leader 中同步数据保持和 leader 数据的同步。leader 发生故障时某个 follower 会成为新的 leader。Kafka/ ZooKeeper用来管理 Producer、broker、Consumer并协调请求和转发。旧版本 Kafka 依赖 ZK 存储集群元数据Broker、分区、offset、控制器信息Kafka 2.8 支持 KRaft 模式移除了 Zookeeperhttps://cloud.tencent.com/developer/article/2109304kafka 关键特性消息持久化消息写入磁盘支持设置保留时间默认 7 天消息消费后不会立刻删除高吞吐顺序写磁盘、零拷贝、批量发送多副本高可用Leader 宕机Follower 自动切换消息的消费点对点、发布/订阅Kafka 支持两种消息传输模型点对点即一对一消费者主动拉取数据消息收到后消息清除。消息生产者生产消息发送到 Queue 中然后消费者从 Queue 中取出并且消费消息。消息被消费以后Queue 中不再有存储所以消费者不可能消费到已经被消费的消息。Queue 支持存在多个消费者但对于一个消息而言只有一个消费者可以消费。发布 / 订阅模式一对多消费者消费数据之后不会清除消息。多个不同消费组消息会被每个消费组各自消费。消息生产者发布将消息发布到 topic 中同时有多个消息消费者订阅消费该消息。和点对点方式不同发布到 topic 中的消息会被所有订阅者消费。易错点 ( 重点 ) 一个 Topic 中的一个分区只能被同一个 Consumer Group 中的一个消费者消费其他消费者不能进行消费。这里的一个消费者指的是运行消费者应用的进程也可以是一个线程。消费者消费完消息后消息不会立刻删除。消费 offset 和 消息保留策略互相独立。消费者有没有消费不影响消息什么时候被删掉。哪怕没人消费假设设置的策略是到 7 天删除即使不到7天消息被消费完但是消息依然存在磁盘。消息不会因为被消费完而提前删除。offset 删除不等价于消息删除。offset 是消费位置存在__consumer_offsets 主题 业务消息存在业务 topic 分区日志两套独立数据。消息保存时间长短与副本因子无关。副本只是多一份备份到期统一清理。 假如副本有 3 份则7 天后 3 份副本全部删除。消息投递语义at most once最多一次消息可能丢失不会重复消费前提交 offsetat least once至少一次消息不会丢失可能重复默认推荐消费成功后提交 offsetexactly once精确一次Kafka Streams 事务实现业务侧一般靠幂等处理重复消息Kafka 消息 / 副本 存储时长副本没有独立的过期删除时间控制数据保留多久的是【消息保留策略】分区所有副本一起遵守该策略。 消息到达保留阈值后分区日志分段统一删除Leader、Follower 副本同步清理。Kafka不是单条消息过期删除而是按「日志段 log segment」为单位删除。 一个 segment 文件写满后关闭新建下一个 segment只有整个 segment 内所有消息都超过保留时间才会删除该文件。可以单独给某个 topic 设置保留时间优先级高于服务端全局配置。副本什么时候删除数据消息写入 LeaderFollower 同步日志段副本数据和 Leader 基本一致后台日志清理线程log cleaner判定旧 segment 达到保留条件Leader 先删除本地日志段之后 Follower 副本也会同步删除对应 segment安装 kafka文档https://kafka.apache.org/43/getting-started/官网https://kafka.apache.org/downloads 下载二进制包并解压路径不要带中文、空格。解压后的 Kafka 目录包含以下重要文件夹bin/windows/- Windows批处理脚本config/- 配置文件libs/- 依赖库logs/- 日志文件启动后生成启动顺序KRaft直接启动 kafkaZK 模式先启动 zookeeper再启动 kafka注意事项关闭服务直接关闭 cmd 窗口即可不要强行反复格式化 kraft 目录如果端口被占用修改server.properties内listenersPLAINTEXT://localhost:9093Windows 下 wsl 运行 kafka 稳定性远高于原生 cmd如果使用 python 连接 kafka连接地址填写localhost:9092KRaft 模式 (无需 Zookeeper)生成集群唯一 ID只需执行一次打开 CMD进入 kafka 根目录进入目录cd D:\kafka\bin\windows\ 执行命令kafka-storage.bat random-uuid 输出类似abcdefg-xxxx复制保存这个 uuid修改配置文件路径\kafka\config\server.properties找到并修改# 填入刚才生成的uuid node.id1 controller.quorum.voters1localhost:9093 listenersPLAINTEXT://localhost:9092 advertised.listenersPLAINTEXT://localhost:9092格式化数据目录只执行第一次不要重复执行执行命令kafka-storage.bat format -t 生成的uuid -c config/kraft/server.properties启动 Kafka 服务.\bin\windows\kafka-server-start.bat config\kraft\server.properties看到日志不再报错代表启动成功不要关闭这个 cmd 窗口停止服务直接关闭运行对应的 cmd 窗口即可。优雅关闭推荐不强制杀进程# kraft模式停止 .\bin\windows\kafka-server-stop.bat config/kraft/server.properties # zk模式停止 .\bin\windows\kafka-server-stop.bat config/server.properties .\bin\windows\zookeeper-server-stop.bat config/zookeeper.properties传统模式 (Zookeeper Kafka)Kafka 教程(图文详解可视化管理工具)https://www.cnblogs.com/ycfenxi/p/19200689启动 Zookeeper。新开 cmd执行下面命令 (默认端口2181窗口保持打开) cd D:\kafka_2.13-3.6.1.\bin\windows\zookeeper-server-start.bat config\zookeeper.properties启动 Kafka新开 cmd:端口默认9092.\bin\windows\kafka-server-start.bat config\server.properties常用 Kafka 命令# 创建topic kafka-topics --create --topic test_topic --bootstrap-server 127.0.0.1:9092 --partitions 3 --replication-factor 1 # 查看topic kafka-topics --list --bootstrap-server 127.0.0.1:9092 # 控制台生产者 kafka-console-producer --topic test_topic --bootstrap-server 127.0.0.1:9092 # 控制台消费者从头消费 kafka-console-consumer --topic test_topic --bootstrap-server 127.0.0.1:9092 --from-beginning创建 topic.\bin\windows\kafka-topics.bat --create --topic test_topic --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1查看 topic 列表.\bin\windows\kafka-topics.bat --list --bootstrap-server localhost:9092控制台生产者发消息.\bin\windows\kafka-console-producer.bat --topic test_topic --bootstrap-server localhost:9092控制台消费者收消息.\bin\windows\kafka-console-consumer.bat --topic test_topic --bootstrap-server localhost:9092 --from-beginning可视化工具Kafka Toolkafka-uihttps://github.com/provectus/kafka-uiKafka Web UIhttps://github.com/obsidiandynamics/kafdropKafka Toolhttps://www.kafkatool.com/download.html以 kafka tool 为例安装完成后配置连接启动 Kafka Tool点击 File → Add New Connection配置连接参数Cluster name: Local KafkaKafka Cluster Version: 3.5Bootstrap servers: localhost:9092点击 Test 测试连接端口占用问题 检查端口占用结束占用进程或修改端口配置netstat -ano | findstr :9092netstat -ano | findstr :2181内存不足问题修改bin\windows\kafka-server-start.bat文件中的JVM参数set KAFKA_HEAP_OPTS-Xmx512M -Xms512M防火墙问题现象远程无法连接Kafka解决方案在Windows防火墙中开放9092和2181端口2、Python 操作 KafkaPython 主流客户端库kafka-python最常用纯 Python 实现confluent-kafkaC 封装性能更高生产环境推荐前置条件本地 / 服务器部署 Kafka启动服务地址127.0.0.1:9092方案1kafka-python简单易上手 pip install kafka-python方案2高性能 confluent-kafkapip install confluent-kafka常见问题有序性问题Kafka 只保证分区内有序Topic 全局无序如果需要全局有序Topic 只能设置1 个分区并发能力大幅下降。消息重复消费使用手动提交 offset 依然可能重复处理完业务提交 offset 前进程宕机 解决方案业务层实现幂等性唯一 key 去重。消费组与分区数量消费者数量不能超过分区数多余消费者空闲想要提高并发需要增加分区。消息丢失场景生产者acks0/1Leader 宕机消息未同步副本。消费者自动提交 offset消息未处理完成就提交 offset消息积压消费速度 生产速度 → 分区消息堆积排查消费逻辑慢、分区过少、消费者数量不足。序列化规范统一使用 json/protobuf不要直接传输字符串客户端编解码保持一致。使用 kafka-pythonfrom kafka import KafkaProducer import json # 初始化生产者 producer KafkaProducer( bootstrap_servers[127.0.0.1:9092], # kafka集群地址多个用逗号分隔 # 序列化value转为bytes value_serializerlambda v: json.dumps(v).encode(utf-8), # key序列化可选 key_serializerlambda k: str(k).encode(utf-8), acks1, # 确认机制0/1/all生产推荐1强一致性选all retries3 # 发送失败重试次数 ) # 1. 发送普通消息异步 msg {name: test, data: hello kafka} # send(主题, value, keyxxx) future producer.send(test_topic, valuemsg) # 等待发送结果同步阻塞可选 try: record_metadata future.get(timeout10) print(f发送成功: 分区{record_metadata.partition}, offset{record_metadata.offset}) except Exception as e: print(发送失败, e) # 2. 带key发送相同key进入同一个分区 producer.send(test_topic, key1001, value{user_id:1001, msg:user message}) # 批量刷新缓冲区程序退出前必须调用 producer.flush()参数说明acks0不等待 broker 响应吞吐最高容易丢消息acks1Leader 写入成功即返回平衡可靠性与性能默认acksall所有副本同步完成才返回可靠性最高性能低消费者 Consumer核心from kafka import KafkaConsumer import json consumer KafkaConsumer( test_topic, # 订阅主题可以多个 [topic1,topic2] bootstrap_servers[127.0.0.1:9092], group_idmy_group_01, # 消费组ID同一个组分摊分区 auto_offset_resetearliest, # auto_offset_reset可选值 # earliest没有offset时从头开始消费 # latest没有offset时只消费新产生消息 enable_auto_commitTrue, # 自动提交offset默认True auto_commit_interval_ms1000, # 自动提交间隔 value_deserializerlambda m: json.loads(m.decode(utf-8)) ) # 持续拉取消息 print(开始监听消息...) for msg in consumer: print(*50) print(f主题: {msg.topic}) print(f分区: {msg.partition}) print(foffset: {msg.offset}) print(fkey: {msg.key}) print(f消息内容: {msg.value})生产重要提醒enable_auto_commitTrue会定时自动提交 offset存在消息丢失风险风险场景offset 提交成功但程序处理消息中途崩溃 → 消息丢失生产推荐关闭自动提交手动提交 offsetat least once手动提交 offset 消费者示例from kafka import KafkaConsumer import json consumer KafkaConsumer( test_topic, bootstrap_servers[127.0.0.1:9092], group_idmy_group_01, auto_offset_resetearliest, enable_auto_commitFalse, # 关闭自动提交 value_deserializerlambda m: json.loads(m.decode(utf-8)) ) for msg in consumer: try: # 执行业务逻辑 print(处理消息, msg.value) # 业务处理成功后手动提交offset consumer.commit() except Exception as e: # 处理失败不提交offset下次重启重新消费这条消息 print(消息处理异常不提交offset, e)常用操作指定起始 offset 消费# 获取分区信息 from kafka import TopicPartition tp TopicPartition(test_topic, partition0) consumer.assign([tp]) consumer.seek(tp, offset10) # 从offset10开始消费使用 confluent-kafkakafka-python 纯 Python 实现高并发场景性能较差生产优先使用confluent-kafka生产者from confluent_kafka import Producer import json conf { bootstrap.servers: 127.0.0.1:9092, acks: 1 } p Producer(conf) # 消息送达回调 def delivery_report(err, msg): if err is not None: print(f消息发送失败: {err}) else: print(f消息发送成功: {msg.topic()} [{msg.partition()}] offset{msg.offset()}) data json.dumps({msg:confluent producer test}).encode() p.produce(test_topic, valuedata, on_deliverydelivery_report) # 轮询触发回调 p.poll(0) # 刷新 p.flush()消费者from confluent_kafka import Consumer, KafkaError import json conf { bootstrap.servers: 127.0.0.1:9092, group.id: confluent_group, auto.offset.reset: earliest, enable.auto.commit: False # 手动提交 } c Consumer(conf) c.subscribe([test_topic]) try: while True: msg c.consume(1, timeout1.0) if msg is None: continue if msg.error(): if msg.error().code() KafkaError._PARTITION_EOF: continue else: raise msg.error() val json.loads(msg.value().decode()) print(收到消息, val) # 业务成功后提交 c.commit(msg) except KeyboardInterrupt: pass finally: c.close()封装 Python Kafka 工具类实现生产者 消费者支持重试、日志、幂等以及 Kafka 事务、批量发送、异步消费采用confluent-kafkalibrdkafka 底层高性能生产首选 功能清单生产者批量发送、事务支持、发送重试、回调日志消费者手动提交 offset、异步消费封装、异常重试、幂等设计样板统一日志、异常捕获、可配置化安装依赖pip install confluent-kafka python-dotenv目录结构kafka_client/├── .env # 配置文件├── kafka_producer.py # 生产者封装批量、事务、重试├── kafka_consumer.py # 消费者封装手动提交、异步消费└── demo.py # 使用示例.env 配置文件# kafka集群地址 KAFKA_BOOTSTRAP_SERVERS127.0.0.1:9092 # 事务ID前缀开启事务必须配置 KAFKA_TRANSACTION_ID_PREFIXtx_ # 消息超时 KAFKA_MESSAGE_TIMEOUT_MS120000 # 批量 linger 等待时间ms KAFKA_LING_MS5 # 批量大小上限 KAFKA_BATCH_SIZE16384kafka_producer.py 生产者封装。支持普通发送、批量发送、Kafka 事务发送、失败回调、重试机制import json import logging from typing import List, Dict, Optional from confluent_kafka import Producer, KafkaError, KafkaException from dotenv import load_dotenv import os load_dotenv() logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s) logger logging.getLogger(KafkaProducer) class KafkaProducerClient: def __init__(self, transactional: bool False): :param transactional: 是否开启事务模式 self.bootstrap_servers os.getenv(KAFKA_BOOTSTRAP_SERVERS) self.ling_ms int(os.getenv(KAFKA_LING_MS, 5)) self.batch_size int(os.getenv(KAFKA_BATCH_SIZE, 16384)) self.transactional transactional conf { bootstrap.servers: self.bootstrap_servers, acks: 1, linger.ms: self.ling_ms, batch.size: self.batch_size, retries: 3, # 发送重试次数 message.timeout.ms: int(os.getenv(KAFKA_MESSAGE_TIMEOUT_MS, 120000)), queue.buffering.max.messages: 100000, } if transactional: # 事务模式必须配置 transactional.id tx_id f{os.getenv(KAFKA_TRANSACTION_ID_PREFIX)}{os.urandom(4).hex()} conf[transactional.id] tx_id logger.info(f开启事务模式transactional.id{tx_id}) self.producer Producer(conf) if self.transactional: self.producer.init_transactions() def _delivery_report(self, err, msg): 消息发送回调 if err is not None: logger.error(f消息发送失败: {err} topic{msg.topic()}) else: logger.debug(f消息发送成功 topic{msg.topic()} partition{msg.partition()} offset{msg.offset()}) def send(self, topic: str, data: Dict, key: Optional[str] None): 单条消息发送异步 payload json.dumps(data, ensure_asciiFalse).encode(utf-8) key_bytes key.encode(utf-8) if key else None try: self.producer.produce( topictopic, valuepayload, keykey_bytes, on_deliveryself._delivery_report ) self.producer.poll(0) except KafkaException as e: logger.exception(fproduce异常 topic{topic}) raise e def batch_send(self, topic: str, msg_list: List[Dict], key_list: Optional[List[str]] None): 批量发送多条消息 key_list key_list or [None] * len(msg_list) for idx, data in enumerate(msg_list): key key_list[idx] payload json.dumps(data, ensure_asciiFalse).encode(utf-8) key_bytes key.encode(utf-8) if key else None self.producer.produce(topic, valuepayload, keykey_bytes, on_deliveryself._delivery_report) self.producer.poll(0) def transaction_send(self, topic: str, msg_list: List[Dict]): Kafka事务发送要么全部成功要么全部失败 适用场景多条消息原子写入 if not self.transactional: raise RuntimeError(实例创建时必须指定 transactionalTrue 才能使用事务) try: self.producer.begin_transaction() for data in msg_list: payload json.dumps(data, ensure_asciiFalse).encode(utf-8) self.producer.produce(topic, valuepayload, on_deliveryself._delivery_report) self.producer.commit_transaction() logger.info(f事务提交成功消息数量:{len(msg_list)}) except KafkaException as e: self.producer.abort_transaction() logger.error(f事务回滚异常:{e}) raise e def flush(self, timeout: int 5000): 等待缓冲区消息全部发送完成程序退出前调用 self.producer.flush(timeouttimeout)kafka_consumer.py 消费者封装。特性如下手动提交 offsetat-least-once支持异步消费线程池消费异常捕获、幂等处理样板支持重启从头 / 最新消费import json import logging import time import threading from concurrent.futures import ThreadPoolExecutor, TimeoutError from typing import Callable, Optional from confluent_kafka import Consumer, KafkaError, KafkaException from dotenv import load_dotenv import os load_dotenv() logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s) logger logging.getLogger(KafkaConsumer) class KafkaConsumerClient: def __init__( self, group_id: str, topics: list, auto_offset_reset: str earliest, enable_auto_commit: bool False ): self.bootstrap_servers os.getenv(KAFKA_BOOTSTRAP_SERVERS) self.group_id group_id self.topics topics self.auto_offset_reset auto_offset_reset self.enable_auto_commit enable_auto_commit conf { bootstrap.servers: self.bootstrap_servers, group.id: self.group_id, auto.offset.reset: self.auto_offset_reset, enable.auto.commit: self.enable_auto_commit, fetch.min.bytes: 1, fetch.max.wait.ms: 500, } self.consumer Consumer(conf) self.consumer.subscribe(self.topics) self.running False def _process_msg(self, msg_handler: Callable, raw_msg): 单条消息处理 try: value json.loads(raw_msg.value().decode(utf-8)) key raw_msg.key().decode(utf-8) if raw_msg.key() else None # 幂等性样板逻辑 # 业务建议每条消息携带唯一biz_id消费前先查询数据库/redis判断是否已处理 # biz_id value.get(biz_id) # if redis.exists(biz_id): # logger.warning(f消息已处理跳过 biz_id{biz_id}) # return # msg_handler(value, key, raw_msg) except Exception as e: logger.exception(f消息处理异常 topic{raw_msg.topic()}) raise e def start_sync_consume(self, msg_handler: Callable, poll_timeout: float 1.0): 同步消费串行简单稳定 self.running True logger.info(f开始同步消费 topics{self.topics} group{self.group_id}) try: while self.running: msg self.consumer.consume(num_messages1, timeoutpoll_timeout) if not msg: continue msg msg[0] if msg.error(): if msg.error().code() KafkaError._PARTITION_EOF: continue raise KafkaException(msg.error()) try: self._process_msg(msg_handler, msg) # 业务处理成功手动提交offset if not self.enable_auto_commit: self.consumer.commit(messagemsg, asynchronousFalse) except Exception: # 处理失败不提交offset下次重启重新消费 time.sleep(1) except KeyboardInterrupt: logger.info(收到停止信号) finally: self.close() def start_async_consume(self, msg_handler: Callable, worker_num: int 4, poll_timeout: float 1.0): 异步消费线程池并发处理消息 注意同一个分区消息会乱序如果业务需要分区有序不能开多线程 self.running True logger.info(f开始异步消费 topics{self.topics} group{self.group_id} workers{worker_num}) executor ThreadPoolExecutor(max_workersworker_num) try: while self.running: msg_list self.consumer.consume(num_messages1, timeoutpoll_timeout) if not msg_list: continue msg msg_list[0] if msg.error(): if msg.error().code() KafkaError._PARTITION_EOF: continue raise KafkaException(msg.error()) # 提交线程池执行 future executor.submit(self._process_msg, msg_handler, msg) try: future.result(timeout10) # 任务执行成功才提交offset if not self.enable_auto_commit: self.consumer.commit(messagemsg, asynchronousFalse) except TimeoutError: logger.error(消息处理超时) future.cancel() except Exception: logger.error(消息处理失败不提交offset) except KeyboardInterrupt: logger.info(收到停止信号) finally: self.running False executor.shutdown(waitTrue) self.close() def close(self): logger.info(关闭kafka消费者) self.consumer.close()demo.py 使用示例from kafka_producer import KafkaProducerClient from kafka_consumer import KafkaConsumerClient # 生产者示例 def demo_producer(): # 1.普通生产者 producer KafkaProducerClient(transactionalFalse) # 单条发送 producer.send(demo_topic, {biz_id: order_001, order_name: 手机订单}, keyorder_001) # 批量发送 batch_data [ {biz_id: order_002, amount: 100}, {biz_id: order_003, amount: 299} ] producer.batch_send(demo_topic, batch_data) # 2.事务生产者原子批量 tx_producer KafkaProducerClient(transactionalTrue) tx_data [ {biz_id: tx_001, type: pay}, {biz_id: tx_001, type: notify} ] tx_producer.transaction_send(demo_topic, tx_data) producer.flush() tx_producer.flush() # 消费者业务处理函数 def message_handler(data: dict, key: str, raw_msg): 业务处理逻辑 logger.info(f收到消息 key{key} data{data}) # 模拟业务 # 幂等逻辑写在这里使用biz_id去重 # 数据库操作 / http调用 # 同步消费 def demo_sync_consumer(): consumer KafkaConsumerClient( group_iddemo_group_01, topics[demo_topic], auto_offset_resetearliest, enable_auto_commitFalse ) consumer.start_sync_consume(message_handler) # 异步多线程消费 def demo_async_consumer(): consumer KafkaConsumerClient( group_iddemo_group_02, topics[demo_topic], auto_offset_resetlatest, enable_auto_commitFalse ) # worker_num 根据分区数量调整不要超过分区总数 consumer.start_async_consume(message_handler, worker_num3) if __name__ __main__: import logging logger logging.getLogger() # demo_producer() # demo_sync_consumer() demo_async_consumer()注意事项Kafka 事务限制事务生产者不要混用普通 send 和 transaction_sendtransactional.id 必须全局唯一重启尽量复用可选事务适合「多条消息原子写入同一个 / 多个 topic」开销更大不要滥用事务不能解决下游业务异常仅保证 producer 侧原子写入批量发送原理linger.ms消息不会立刻发送等待一小段时间凑更多消息打包发送提升吞吐高吞吐场景linger.ms5~20低延迟场景linger.ms0异步消费风险异步多线程消费会破坏分区内有序性如果业务要求同一个 key 消息有序禁止异步多线程消费只能同步串行并发上限线程数 ≤ topic 分区数量再多线程无法提升速度offset 提交策略enable_auto_commitTrue极易丢失消息生产禁止使用工具类默认手动提交业务完全成功后再 commit保证至少一次消息积压排查消费逻辑阻塞IO 慢、无超时分区数量太少无法扩容消费者消费者频繁重启offset 反复回滚