kafka:https://github.com/apache/kafka
kafka-python:https://github.com/dpkp/kafka-python
1、Kafka 简介
什么是 kafka
Apache Kafka 是一款分布式、高吞吐、可持久化、分区、多副本的消息队列(消息中间件),基于发布 / 订阅模式。 典型场景:
- 日志收集、实时数据流
- 系统解耦、异步通信
- 流量削峰填谷
- 实时计算(Flink/Spark Streaming 数据源)
Kafka 为什么这么强? 高吞吐:磁盘顺序读写 + 批量压缩。“Kafka 基于磁盘存储,为什么能比内存队列更快?” 答案在于磁盘顺序读写和批量压缩。
- 传统磁盘 IO 慢的核心是 “随机读写”(磁头频繁移动),而 Kafka 中每个 Partition 的数据是按 “日志文件” 形式顺序存储的 ——Producer 写入时只能在文件末尾追加(Append Only),Consumer 读取时也按顺序从头往后读,避免了随机 IO 的开销,磁盘顺序读写的速度甚至接近内存。
- 同时,Kafka 支持对消息进行批量压缩(如 Gzip、Snappy),Producer 将多个小消息打包压缩后发送,减少网络传输量;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 模式,移除了 Zookeeper:https://cloud.tencent.com/developer/article/2109304
kafka 关键特性
- 消息持久化:消息写入磁盘,支持设置保留时间(默认 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(最多一次):消息可能丢失,不会重复(消费前提交 offset)
- at least once(至少一次):消息不会丢失,可能重复(默认推荐,消费成功后提交 offset)
- exactly once(精确一次):Kafka Streams 事务实现,业务侧一般靠幂等处理重复消息
Kafka 消息 / 副本 存储时长
副本没有独立的过期删除时间;控制数据保留多久的是【消息保留策略】,分区所有副本一起遵守该策略。 消息到达保留阈值后,分区日志分段统一删除,Leader、Follower 副本同步清理。
Kafka不是单条消息过期删除,而是按「日志段 log segment」为单位删除。 一个 segment 文件写满后关闭,新建下一个 segment;只有整个 segment 内所有消息都超过保留时间,才会删除该文件。
可以单独给某个 topic 设置保留时间,优先级高于服务端全局配置。
副本什么时候删除数据?
- 消息写入 Leader
- Follower 同步日志段,副本数据和 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:直接启动 kafka
- ZK 模式:先启动 zookeeper,再启动 kafka
注意事项
- 关闭服务:直接关闭 cmd 窗口即可;不要强行反复格式化 kraft 目录
- 如果端口被占用:修改
server.properties内listeners=PLAINTEXT://localhost:9093 - Windows 下 wsl 运行 kafka 稳定性远高于原生 cmd
- 如果使用 python 连接 kafka,连接地址填写:
localhost:9092
KRaft 模式 (无需 Zookeeper)
生成集群唯一 ID(只需执行一次),打开 CMD,进入 kafka 根目录:
进入目录:cd D:\kafka\bin\windows\ 执行命令:kafka-storage.bat random-uuid 输出类似:abcdefg-xxxx,复制保存这个 uuid修改配置文件,路径:\kafka\config\server.properties找到并修改:
# 填入刚才生成的uuid node.id=1 controller.quorum.voters=1@localhost:9093 listeners=PLAINTEXT://localhost:9092 advertised.listeners=PLAINTEXT://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 Tool
kafka-ui:https://github.com/provectus/kafka-ui
Kafka Web UI:https://github.com/obsidiandynamics/kafdrop
Kafka Tool:https://www.kafkatool.com/download.html
以 kafka tool 为例,安装完成后,配置连接:
- 启动 Kafka Tool,点击 File → Add New Connection
- 配置连接参数:
Cluster name: Local Kafka
Kafka Cluster Version: 3.5
Bootstrap servers: localhost:9092 - 点击 Test 测试连接
端口占用问题, 检查端口占用,结束占用进程或修改端口配置
netstat -ano | findstr :9092
netstat -ano | findstr :2181
内存不足问题,修改bin\windows\kafka-server-start.bat文件中的JVM参数:
set "KAFKA_HEAP_OPTS=-Xmx512M -Xms512M"
防火墙问题
现象:远程无法连接Kafka
解决方案::在Windows防火墙中开放9092和2181端口
2、Python 操作 Kafka
Python 主流客户端库:
- kafka-python(最常用,纯 Python 实现)
- confluent-kafka(C 封装,性能更高,生产环境推荐)
前置条件:本地 / 服务器部署 Kafka,启动服务,地址:
127.0.0.1:9092方案1:kafka-python(简单易上手) :pip install kafka-python
方案2:高性能 confluent-kafka:pip install confluent-kafka
常见问题
- 有序性问题:Kafka 只保证分区内有序,Topic 全局无序;如果需要全局有序,Topic 只能设置1 个分区,并发能力大幅下降。
- 消息重复消费:使用手动提交 offset 依然可能重复(处理完业务,提交 offset 前进程宕机); 解决方案:业务层实现幂等性(唯一 key 去重)。
- 消费组与分区数量消费者数量不能超过分区数,多余消费者空闲;想要提高并发,需要增加分区。
- 消息丢失场景:生产者:acks=0/1,Leader 宕机消息未同步副本。消费者:自动提交 offset,消息未处理完成就提交 offset
- 消息积压消费速度 < 生产速度 → 分区消息堆积;排查:消费逻辑慢、分区过少、消费者数量不足。
- 序列化规范统一使用 json/protobuf,不要直接传输字符串,客户端编解码保持一致。
使用 kafka-python
from kafka import KafkaProducer import json # 初始化生产者 producer = KafkaProducer( bootstrap_servers=['127.0.0.1:9092'], # kafka集群地址,多个用逗号分隔 # 序列化:value转为bytes value_serializer=lambda v: json.dumps(v).encode('utf-8'), # key序列化(可选) key_serializer=lambda k: str(k).encode('utf-8'), acks=1, # 确认机制:0/1/all;生产推荐1,强一致性选all retries=3 # 发送失败重试次数 ) # 1. 发送普通消息(异步) msg = {"name": "test", "data": "hello kafka"} # send(主题, value, key=xxx) future = producer.send("test_topic", value=msg) # 等待发送结果(同步阻塞,可选) try: record_metadata = future.get(timeout=10) print(f"发送成功: 分区={record_metadata.partition}, offset={record_metadata.offset}") except Exception as e: print("发送失败", e) # 2. 带key发送(相同key进入同一个分区) producer.send("test_topic", key=1001, value={"user_id":1001, "msg":"user message"}) # 批量刷新缓冲区,程序退出前必须调用 producer.flush()参数说明:
acks=0:不等待 broker 响应,吞吐最高,容易丢消息acks=1:Leader 写入成功即返回(平衡可靠性与性能,默认)acks=all:所有副本同步完成才返回,可靠性最高,性能低
消费者 Consumer(核心)
from kafka import KafkaConsumer import json consumer = KafkaConsumer( "test_topic", # 订阅主题,可以多个 ["topic1","topic2"] bootstrap_servers=['127.0.0.1:9092'], group_id="my_group_01", # 消费组ID,同一个组分摊分区 auto_offset_reset="earliest", # auto_offset_reset可选值: # earliest:没有offset时,从头开始消费 # latest:没有offset时,只消费新产生消息 enable_auto_commit=True, # 自动提交offset(默认True) auto_commit_interval_ms=1000, # 自动提交间隔 value_deserializer=lambda m: json.loads(m.decode('utf-8')) ) # 持续拉取消息 print("开始监听消息...") for msg in consumer: print("="*50) print(f"主题: {msg.topic}") print(f"分区: {msg.partition}") print(f"offset: {msg.offset}") print(f"key: {msg.key}") print(f"消息内容: {msg.value}")生产重要提醒:enable_auto_commit=True会定时自动提交 offset,存在消息丢失风险!
- 风险场景:offset 提交成功,但程序处理消息中途崩溃 → 消息丢失
- 生产推荐:关闭自动提交,手动提交 offset(at least once)
手动提交 offset 消费者示例
from kafka import KafkaConsumer import json consumer = KafkaConsumer( "test_topic", bootstrap_servers=['127.0.0.1:9092'], group_id="my_group_01", auto_offset_reset="earliest", enable_auto_commit=False, # 关闭自动提交! value_deserializer=lambda 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", partition=0) consumer.assign([tp]) consumer.seek(tp, offset=10) # 从offset=10开始消费使用 confluent-kafka
kafka-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', value=data, on_delivery=delivery_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, timeout=1.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-kafka(librdkafka 底层,高性能,生产首选) 功能清单:
- 生产者:批量发送、事务支持、发送重试、回调日志
- 消费者:手动提交 offset、异步消费封装、异常重试、幂等设计样板
- 统一日志、异常捕获、可配置化
安装依赖:pip install confluent-kafka python-dotenv
目录结构
kafka_client/
├── .env # 配置文件
├── kafka_producer.py # 生产者封装(批量、事务、重试)
├── kafka_consumer.py # 消费者封装(手动提交、异步消费)
└── demo.py # 使用示例
.env 配置文件
# kafka集群地址 KAFKA_BOOTSTRAP_SERVERS=127.0.0.1:9092 # 事务ID前缀(开启事务必须配置) KAFKA_TRANSACTION_ID_PREFIX=tx_ # 消息超时 KAFKA_MESSAGE_TIMEOUT_MS=120000 # 批量 linger 等待时间,ms KAFKA_LING_MS=5 # 批量大小上限 KAFKA_BATCH_SIZE=16384kafka_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(level=logging.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_ascii=False).encode("utf-8") key_bytes = key.encode("utf-8") if key else None try: self.producer.produce( topic=topic, value=payload, key=key_bytes, on_delivery=self._delivery_report ) self.producer.poll(0) except KafkaException as e: logger.exception(f"produce异常 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_ascii=False).encode("utf-8") key_bytes = key.encode("utf-8") if key else None self.producer.produce(topic, value=payload, key=key_bytes, on_delivery=self._delivery_report) self.producer.poll(0) def transaction_send(self, topic: str, msg_list: List[Dict]): """ Kafka事务发送:要么全部成功,要么全部失败 适用场景:多条消息原子写入 """ if not self.transactional: raise RuntimeError("实例创建时必须指定 transactional=True 才能使用事务") try: self.producer.begin_transaction() for data in msg_list: payload = json.dumps(data, ensure_ascii=False).encode("utf-8") self.producer.produce(topic, value=payload, on_delivery=self._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(timeout=timeout)kafka_consumer.py 消费者封装。特性如下:
- 手动提交 offset(at-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(level=logging.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_messages=1, timeout=poll_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(message=msg, asynchronous=False) 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_workers=worker_num) try: while self.running: msg_list = self.consumer.consume(num_messages=1, timeout=poll_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(timeout=10) # 任务执行成功才提交offset if not self.enable_auto_commit: self.consumer.commit(message=msg, asynchronous=False) except TimeoutError: logger.error("消息处理超时") future.cancel() except Exception: logger.error("消息处理失败,不提交offset") except KeyboardInterrupt: logger.info("收到停止信号") finally: self.running = False executor.shutdown(wait=True) 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(transactional=False) # 单条发送 producer.send("demo_topic", {"biz_id": "order_001", "order_name": "手机订单"}, key="order_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(transactional=True) 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_id="demo_group_01", topics=["demo_topic"], auto_offset_reset="earliest", enable_auto_commit=False ) consumer.start_sync_consume(message_handler) # ====================== 异步多线程消费 ====================== def demo_async_consumer(): consumer = KafkaConsumerClient( group_id="demo_group_02", topics=["demo_topic"], auto_offset_reset="latest", enable_auto_commit=False ) # worker_num 根据分区数量调整,不要超过分区总数 consumer.start_async_consume(message_handler, worker_num=3) if __name__ == "__main__": import logging logger = logging.getLogger() # demo_producer() # demo_sync_consumer() demo_async_consumer()注意事项
Kafka 事务限制
- 事务生产者不要混用普通 send 和 transaction_send
- transactional.id 必须全局唯一,重启尽量复用(可选)
- 事务适合「多条消息原子写入同一个 / 多个 topic」,开销更大,不要滥用
- 事务不能解决下游业务异常,仅保证 producer 侧原子写入
批量发送原理
linger.ms:消息不会立刻发送,等待一小段时间凑更多消息打包发送,提升吞吐;
- 高吞吐场景:
linger.ms=5~20 - 低延迟场景:
linger.ms=0
异步消费风险,异步多线程消费会破坏分区内有序性!
- 如果业务要求同一个 key 消息有序:禁止异步多线程消费,只能同步串行
- 并发上限:线程数 ≤ topic 分区数量,再多线程无法提升速度
offset 提交策略
enable_auto_commit=True:极易丢失消息,生产禁止使用- 工具类默认手动提交:业务完全成功后再 commit,保证至少一次
消息积压排查
- 消费逻辑阻塞(IO 慢、无超时)
- 分区数量太少,无法扩容消费者
- 消费者频繁重启,offset 反复回滚