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}
原因分析
- 元数据重试机制:Kafka Producer发送消息前会先获取主题元数据,当主题不存在时,默认会持续重试元数据请求(
metadata.max.age.ms默认5分钟,metadata.max.retries默认无限制),此过程中不会将异常传递到发送回调。 - 自动创建主题默认开启:Kafka集群默认开启
auto.create.topics.enable=true,Producer发送消息时会自动创建不存在的主题,不会触发InvalidTopicException。 - 异常包装问题:即使元数据请求失败,异常可能被包装在内部执行逻辑中,未直接传递到
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。
验证步骤
- 修改Producer配置并重启程序
- 确认Kafka集群已关闭自动创建主题
- 删除目标主题后发送消息,此时回调应触发
InvalidTopicException处理逻辑,程序正常终止
内容的提问来源于stack exchange,提问作者user22614887
相关产品推荐
相关产品推荐

