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

Kafka生产者/消费者异常:消费者无法接收全部消息排查求助

Kafka生产者快速发送消息后消费者接收不全的问题

问题描述

我正在测试Kafka以评估其在未来项目中的适用性,基于kafka-python库编写了两个Python脚本。当前遇到问题:生产者发送100条消息后,消费者接收的消息数量不足100条。

环境及操作细节

  • 使用默认配置的Kafka Broker
  • 仅启动一个消费者,订阅quickstart-events主题以接收所有消息
  • 生产者控制台会输出已发送的msg1至msg99等消息
  • 消费者仅能接收部分消息(如示例中仅收到msg1至msg92)
  • 关键现象:若在生产者中添加2ms的time.sleep(2/1000)延迟,消费者即可接收全部100条消息。我怀疑是否是Broker无法处理生产者快速发送的消息,但对此存疑——该主题仅创建了一个分区。

生产者代码

# producer.py
from kafka import KafkaProducer
from kafka.errors import KafkaError
import time

producer = KafkaProducer(bootstrap_servers=['msg1:9092'])

# Asynchronous by default
future = producer.send('quickstart-events', b'raw_bytes')

for i in range(1,100):
    key = f"key{i}"
    msg = f"msg{i}"
    print(f"Publish message {msg}")
    future = producer.send(topic='quickstart-events', key=key.encode(), value=msg.encode())
    #putting in the 2ms sleep DOES allow for all messages to be received    
    #time.sleep(2/1000)

消费者代码

# consumer.py
from kafka import KafkaConsumer

# To consume latest messages and auto-commit offsets
consumer = KafkaConsumer('quickstart-events',
                     group_id='my-group',
                     bootstrap_servers=['msg1:9092'])

for message in consumer:
    # message value and key are raw bytes -- decode if necessary!
    # e.g., for unicode: `message.value.decode('utf-8')`
    print ("%s:%d:%d: key=%s value=%s" % (message.topic, message.partition,
                                      message.offset, message.key,
                                      message.value))

生产者输出

Publish message msg1
Publish message msg2
Publish message msg3
....
Publish message msg97
Publish message msg98
Publish message msg99

消费者输出

quickstart-events:0:3850: key=None value=b'raw_bytes'
quickstart-events:0:3851: key=b'key1' value=b'msg1'
quickstart-events:0:3852: key=b'key2' value=b'msg2'
quickstart-events:0:3853: key=b'key3' value=b'msg3'
.....
quickstart-events:0:3940: key=b'key90' value=b'msg90'
quickstart-events:0:3941: key=b'key91' value=b'msg91'
quickstart-events:0:3942: key=b'key92' value=b'msg92'

问题原因

这不是Broker处理能力的问题,核心在于kafka-python生产者的异步发送机制+脚本未等待消息完全发送。

producer.send()默认是异步操作,消息会先存入本地发送缓冲区,由后台线程批量发送到Broker。当脚本快速循环完100条消息后,可能直接退出,此时缓冲区中剩余的消息还没来得及发送到Broker,导致消费者收不到。

添加time.sleep(2/1000)后,相当于给后台线程留出了足够的时间把缓冲区消息发送完毕,所以所有消息都能被Broker接收。

解决方案

1. 强制刷新缓冲区

在生产者脚本循环结束后,调用producer.flush(),它会阻塞直到所有缓冲区中的消息都发送到Broker并得到确认:

# 在for循环结束后添加
producer.flush()

2. 等待单条消息发送确认

如果需要确保每条消息都发送成功,可以等待send()返回的Future对象完成,适合对消息可靠性要求极高的场景:

for i in range(1,100):
    key = f"key{i}"
    msg = f"msg{i}"
    print(f"Publish message {msg}")
    future = producer.send(topic='quickstart-events', key=key.encode(), value=msg.encode())
    # 等待消息发送完成,超时时间10秒
    try:
        record_metadata = future.get(timeout=10)
    except KafkaError:
        # 这里可以添加发送失败的处理逻辑
        pass

3. 优雅关闭生产者

在脚本末尾调用producer.close(),它会自动执行flush()并关闭后台线程,确保所有消息都被发送:

# 脚本最后添加
producer.close()

4. 调整批量发送参数

修改生产者的linger_ms参数,让生产者等待一小段时间攒够批量再发送,既提升效率也减少消息残留:

producer = KafkaProducer(
    bootstrap_servers=['msg1:9092'],
    linger_ms=5  # 等待5毫秒再发送批量消息
)

验证方法

修改生产者脚本后重新运行,消费者应该能收到全部100条消息。也可以用Kafka自带的kafka-consumer-groups.sh工具查看主题偏移量,确认Broker已接收所有消息,排除消费者端的问题。


内容的提问来源于stack exchange,提问作者Bik Dhaliwal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 10:34:55