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

为何调用Kafka Producer的producer.close()方法会阻塞数分钟?

问题根因初步判断

你遇到的阻塞时长刚好匹配Kafka生产者默认的delivery.timeout.ms参数值(默认120000ms即2分钟),本质是生产者关闭时等待本地缓冲区未发送完成的消息最终返回结果(要么发送成功、要么重试超时失败),说明你生产的消息大概率一直发送失败,持续重试直到超时才结束,所以卡在producer.close()步骤。

排查步骤

  • 首先确认topic配置合法性
    你贴出的topic信息存在明显异常:Replicas: 0、Isr: 0,正常Kafka分区的副本数最小为1,副本列表为空、ISR为空的情况下,分区无法正常写入数据,broker会一直无法返回合法的ack给生产者。
    执行以下命令重新校验topic配置:
    kafka-topics.sh --describe --bootstrap-server <你的broker地址> --topic myTopic
    如果确认副本数为0,删除异常topic后重新创建,设置合理的副本数,例如3副本集群设置--replication-factor 3。

  • 检查生产者核心配置
    你的代码中未显式配置以下关键参数,默认配置极易导致长时间阻塞:

    1. retries:默认值为Integer.MAX_VALUE,发送失败时会无限重试直到超过投递超时时间
    2. retry.backoff.ms:默认值100ms,两次重试的间隔时间
    3. delivery.timeout.ms:默认值120000ms,也就是你遇到的2分钟超时阈值
      可以临时在生产者配置中添加props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 3000);,如果修改后阻塞时间变成3秒,即可验证是重试超时导致的问题。
  • 添加发送回调查看具体错误
    producer.send()是异步方法,你循环1秒跑完只是把消息写入了本地缓冲区,并未实际完成发送。给send方法添加回调即可直观看到发送错误信息:

producer.send(new ProducerRecord<>(AppConfigs.topicName, i, "msg" + i), (metadata, exception) -> {
    if (exception != null) {
        logger.error("消息发送失败", exception);
    }
});
  • 校验broker运行状态
    查看3台broker的server.log日志,确认是否存在分区不可用、ISR异常收缩、节点间网络不通的报错,同时确认生产者配置的bootstrap.servers地址正确,客户端可以正常访问所有broker的服务端口。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 22:06:03