You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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. 关闭自动提交偏移量:确保只有消息处理完成后才提交偏移量,避免重复消费或漏消费。
  2. 设置单次拉取消息数为1:让消费者每次只获取一条消息,处理完成后再拉取下一条。
  3. 手动提交偏移量:消息处理完成后手动提交,确认消息已被处理。

修改后的消费者代码如下:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 04:22:07