Confluent Kafka Python生产者设置ACKS=all/2时无法生产消息
咱们先拆解核心矛盾:你的Topic设置了min.insync.replicas=2,集群有3个Broker,当生产者把acks设为2或all时消息发不出去,但acks=1正常,而且本地运行没报错——这其实是因为你没捕获到异步发送的错误,再结合之前Lambda遇到的NOT_ENOUGH_REPLICAS错误,大概率是同步副本(ISR)数量不足导致的。
第一步:排查Topic的同步副本(ISR)状态
先确认你的test Topic当前的同步副本数量是否满足min.insync.replicas=2的要求。在终端执行Kafka的topic描述命令:
kafka-topics.sh --describe --topic test --bootstrap-server localhost:9092
重点看输出里的Isr字段,比如如果输出是Isr: 1,说明只有1个Broker在同步副本集合里,这时候当acks=2或all时,生产者需要至少2个副本确认,但根本找不到足够的同步副本,自然无法完成消息写入。
如果ISR数量不足,你需要检查:
- 是否有Broker宕机或者未正常启动?
- Broker之间的网络是否连通?(比如防火墙、端口映射问题导致副本同步失败)
- Broker日志里有没有
UnderReplicatedPartitions或者副本同步相关的错误日志。
第二步:修复生产者代码的错误捕获逻辑
你的现有代码有个关键问题:confluent-kafka的生产者是异步发送的,produce方法不会立即抛出异常,错误会在flush或者通过回调函数返回。你没设置回调,也没检查flush的返回值,所以即使发送失败,也看不到错误信息。
修改后的代码应该加上投递回调和超时配置,同时检查flush的结果:
from confluent_kafka import Producer, KafkaError import logging def get_logger(): logger = logging.getLogger(__name__) logger.setLevel(logging.ERROR) # 添加控制台输出,方便本地调试 handler = logging.StreamHandler() logger.addHandler(handler) return logger def get_config(): conf = { 'bootstrap.servers': 'localhost:9092', 'acks': '2', # 设置请求超时,避免无限等待 'request.timeout.ms': 10000, 'delivery.timeout.ms': 15000, # 全局错误回调,捕获连接、认证等全局错误 'error_cb': lambda err: get_logger().error(f"Kafka全局错误: {err}") } return conf def get_producer_config(): return Producer(get_config()) def delivery_report(err, msg): """投递结果回调函数,捕获单条消息的发送错误""" if err is not None: get_logger().error(f"消息投递失败: {err}") else: get_logger().info(f"消息已投递到 {msg.topic()} [{msg.partition()}]") try: producer = get_producer_config() # 发送消息时指定回调函数 producer.produce('test', 'test message from local app', callback=delivery_report) # 调用flush并指定超时时间,返回值是未成功发送的消息数量 unflushed_count = producer.flush(timeout=10) if unflushed_count > 0: get_logger().error(f"有 {unflushed_count} 条消息未能成功刷新到Kafka") except KafkaError as error: get_logger().error(f"生产者初始化或全局错误: {str(error)}")
这样修改后,你就能清楚看到消息发送失败的具体原因,比如是否真的是NOT_ENOUGH_REPLICAS。
第三步:确认Topic的副本配置
另外,别忘了检查Topic的副本数:如果你的test Topic副本数设置为1,那即使集群有3个Broker,min.insync.replicas=2也完全不合理——因为最多只有1个副本存在,根本达不到2个同步副本的要求。
同样用kafka-topics.sh --describe命令看ReplicationFactor字段,如果副本数小于2,你需要修改Topic的副本数:
kafka-topics.sh --alter --topic test --replication-factor 3 --bootstrap-server localhost:9092
(修改副本数需要Broker集群正常运行,否则可能失败)
最后:验证网络连通性
确保你的本地生产者能连接到所有3个Broker,而不仅仅是localhost:9092。比如如果Broker的advertised.listeners配置的是其他地址,你需要在生产者的bootstrap.servers里填写所有Broker的地址(比如broker1:9092,broker2:9092,broker3:9092),否则生产者可能无法和其他Broker通信,导致无法获取足够的副本确认。
内容的提问来源于stack exchange,提问作者AnonymousAlias

