Python实现Kafka生产者与4个消费者的动态任务负载均衡
实现Kafka生产者与4个消费者的动态负载均衡
问题说明
已创建包含4个分区的numbers主题,当前Kafka默认按分区平均分配消息给消费者组内的4个消费者。需要调整为消费者完成当前任务后再分配新消息的模式,以此提升整体处理效率。
已创建主题的Bash命令
kafka-topics --bootstrap-server localhost:9092 --create --topic numbers --partitions 4 --replication-factor 1
现有代码
生产者代码
from time import sleep from json import dumps from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers=['localhost:9092'], value_serializer=lambda x: dumps(x).encode('utf-8')) for e in range(100): data = {'number' : e} producer.send('numbers', value=data) print(f"Sending data : {data}") sleep(5)
消费者代码
import json, time from kafka import KafkaConsumer print("Connecting to consumer ...") consumer = KafkaConsumer( 'numbers', bootstrap_servers=['localhost:9092'], auto_offset_reset='earliest', enable_auto_commit=True, group_id='my-group', value_deserializer=lambda x: json.loads(x.decode('utf-8'))) for message in consumer: print(f"{message.value}") time.sleep(1)
解决方案
Kafka默认的消费者分配机制是基于分区的静态分配,每个分区只能被消费者组内的一个消费者消费。要实现“处理完再分配”的动态负载,需调整消费者的配置与消费逻辑:
- 关闭自动提交偏移量:确保只有消息处理完成后才提交偏移量,避免重复消费或漏消费。
- 设置单次拉取消息数为1:让消费者每次只获取一条消息,处理完成后再拉取下一条。
- 手动提交偏移量:消息处理完成后手动提交,确认消息已被处理。
修改后的消费者代码如下:
import json, time from kafka import KafkaConsumer print("Connecting to consumer ...") consumer = KafkaConsumer( 'numbers', bootstrap_servers=['localhost:9092'], auto_offset_reset='earliest', enable_auto_commit=False, # 关闭自动提交 group_id='my-group', max_poll_records=1, # 单次拉取1条消息 value_deserializer=lambda x: json.loads(x.decode('utf-8'))) try: while True: # 拉取消息,超时时间设为1秒 messages = consumer.poll(timeout_ms=1000) for topic_partition, records in messages.items(): for record in records: print(f"{record.value}") time.sleep(1) # 模拟业务处理耗时 # 手动提交当前消息的偏移量 consumer.commit({topic_partition: record.offset + 1}) finally: consumer.close()
原理说明
- 保持4个分区与4个消费者的对应关系(每个消费者分配一个分区),保留并行处理的能力。
- 每个消费者每次仅处理一条消息,处理完成后立即提交偏移量并拉取下一条,处理速度快的消费者会更快完成自身分区内的消息处理,不会被处理慢的消费者拖慢整体进度。
- 若需要更灵活的跨分区负载均衡,可将主题设为单分区,但这会牺牲并行性,仅适合消息量不大的场景。
内容的提问来源于stack exchange,提问作者Ali Esmaeili
相关产品推荐
相关产品推荐

