【发布时间】:2021-04-02 17:37:20
【问题描述】:
我对 Kafka 比较陌生,我正在尝试在主题上发送消息后生成消费者。
单个生产者在不同的分区上发送 200 条消息。
consumer-1 已经在运行,consumer-1 正在监听所有 4 个分区,并且在几秒钟内,同一组中的另一个消费者 (consumer-2) 启动了。然后 Kafka 触发组的重新平衡,但所有最初的 200 条 msg 都将发送给 consumer-1,而新的 msg 将发送到 consumer-2。
新建消费者后,描述消费者组API展示Warning: Consumer group 'ner_group' is rebalancing.
谁能建议如何在消费者组中添加新消费者,以读取已经发送到某个主题的消息,而另一个消费者已经在阅读该主题?
下面是消费者和生产者的配置。
制片人
producer = KafkaProducer(bootstrap_servers=[f'{ip}:9092'],
value_serializer=lambda x:
dumps(x).encode('utf-8'))
for e in range(200):
data = {'count':e ,'time':time.time() }
producer.send('ner', value=data)
消费者
consumer = KafkaConsumer(
'ner',
bootstrap_servers=[f'{ip}:9092'],
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='ner_group',
value_deserializer=lambda x: loads(x.decode('utf-8'))
)
for message in consumer:
msg = message.value
time.sleep(1)
我正在多次运行消费者脚本。
【问题讨论】:
-
循环后需要 producer.flush()
标签: python apache-kafka kafka-consumer-api