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

Kafka Producer.send回调在消息发送失败时未执行的原因排查

Kafka异步发送时InvalidTopicException未触发回调异常处理的问题分析与解决

问题背景

使用org.apache.kafka.clients客户端异步发送Kafka消息,依赖Callback处理发送结果,设计目标是遇到InvalidTopicException等不可恢复错误时终止程序。正常发送时回调能正确打印元数据,但删除主题模拟异常时,仅输出NetworkClient的警告日志,回调内的异常处理逻辑未执行。

当前代码实现

消息发送逻辑

try {
    producer.send(record, new Callback() {
        public void onCompletion(RecordMetadata metadata, Exception e) {
            System.out.println(metadata);
            System.out.println(e);
            if (e != null) {
                System.out.println(e);
                try {
                    throw(e);
                } catch (InvalidTopicException er) {
                    logger.error("Unrecoverable error: ", er);
                    System.exit(-1);
                } //...otherfatalexceptions...
                catch (UnknownServerException er) {
                    logger.error("Unrecoverable error: ", er);
                    System.exit(-1);
                } catch (Exception er) {
                    logger.error("Failed to send message to topic:" + topic, er);
                }
            }
        }
    });
} catch (Exception e) {
    logger.error("Failed to send message to topic:" + topic, e);
}

Producer初始化与配置

// 初始化
producer = new KafkaProducer<String, String>(properties);

// 配置项
Properties mainKafkaProperties = new Properties();
mainKafkaProperties.put("bootstrap.servers", kafkaUrl);
mainKafkaProperties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
mainKafkaProperties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
mainKafkaProperties.put("acks", "all");

异常现象

删除主题后发送消息,仅输出警告日志,回调未触发:

[kafka-producer-network-thread | producer-1] WARN org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-1] Error while fetching metadata with correlation id 97 : {em_om_messages_local=UNKNOWN_TOPIC_OR_PARTITION}

原因分析

  1. 元数据重试机制:Kafka Producer发送消息前会先获取主题元数据,当主题不存在时,默认会持续重试元数据请求(metadata.max.age.ms默认5分钟,metadata.max.retries默认无限制),此过程中不会将异常传递到发送回调。
  2. 自动创建主题默认开启:Kafka集群默认开启auto.create.topics.enable=true,Producer发送消息时会自动创建不存在的主题,不会触发InvalidTopicException。
  3. 异常包装问题:即使元数据请求失败,异常可能被包装在内部执行逻辑中,未直接传递到Callback的Exception参数。

解决方案

1. 调整Producer元数据相关配置

添加以下配置,缩短元数据重试周期,让Producer更快将异常传递到回调:

// 缩短元数据过期时间,加快主题不存在的检测速度
mainKafkaProperties.put("metadata.max.age.ms", "5000");
// 设置元数据请求最大重试次数,0表示不重试
mainKafkaProperties.put("metadata.max.retries", "0");
// 设置请求超时时间,避免长时间阻塞
mainKafkaProperties.put("request.timeout.ms", "3000");

2. 优化回调异常处理逻辑

直接判断异常类型(含包装异常),无需抛出再捕获,简化逻辑:

producer.send(record, new Callback() {
    public void onCompletion(RecordMetadata metadata, Exception e) {
        if (e != null) {
            Throwable rootCause = e;
            // 遍历获取根异常
            while (rootCause.getCause() != null) {
                rootCause = rootCause.getCause();
            }
            
            if (rootCause instanceof InvalidTopicException) {
                logger.error("Unrecoverable error: Invalid topic", rootCause);
                System.exit(-1);
            } else if (rootCause instanceof UnknownServerException) {
                logger.error("Unrecoverable error: Unknown server", rootCause);
                System.exit(-1);
            } else {
                logger.error("Failed to send message to topic: " + topic, e);
            }
        } else {
            System.out.println(metadata);
        }
    }
});

3. 关闭Kafka集群自动创建主题

在Kafka broker配置文件中设置:

auto.create.topics.enable=false

确保删除主题后,Producer无法自动创建主题,从而触发InvalidTopicException。

验证步骤

  1. 修改Producer配置并重启程序
  2. 确认Kafka集群已关闭自动创建主题
  3. 删除目标主题后发送消息,此时回调应触发InvalidTopicException处理逻辑,程序正常终止

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 09:41:01