在DLT处理中如何获取@KafkaListener方法抛出的原始异常?
在@DltHandler中获取@KafkaListener抛出的原始异常
可以获取到@KafkaListener方法抛出的原始异常,Spring Kafka在将消息转发到死信队列(DLT)时,会自动把原始异常的关键信息存入消息头中,你可以通过以下两种方式获取:
方式一:直接通过@Header注解提取异常信息
修改DLT处理器方法,添加异常相关的Header参数:
import org.springframework.kafka.annotation.DltHandler; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.annotation.RetryableTopic; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.handler.annotation.Header; import java.nio.charset.StandardCharsets; @KafkaListener(topics = "message", groupId = "1") @RetryableTopic(attempts = "3", backoff = @Backoff(value = 3000L)) public void readMessages(String message) { // 抛出示例异常 throw new IllegalArgumentException("无效的消息格式"); } @DltHandler public void processDltMessages(String message, @Header(KafkaHeaders.EXCEPTION_STACKTRACE) byte[] exceptionStackTrace, @Header(KafkaHeaders.EXCEPTION_MESSAGE) String exceptionMessage) { // 将字节数组转换为可读的栈跟踪字符串 String stackTrace = new String(exceptionStackTrace, StandardCharsets.UTF_8); // 根据异常信息执行自定义逻辑 if (exceptionMessage.contains("无效的消息格式")) { // 处理格式错误的消息 System.out.println("处理格式错误的DLT消息:" + message); } // 打印异常详情 System.out.println("原始异常消息:" + exceptionMessage); System.out.println("原始异常栈:" + stackTrace); }
方式二:通过ConsumerRecord获取所有消息头
如果需要更灵活地操作消息头,可以直接接收ConsumerRecord对象:
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.DltHandler; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.annotation.RetryableTopic; import org.springframework.kafka.support.KafkaHeaders; import java.nio.charset.StandardCharsets; @KafkaListener(topics = "message", groupId = "1") @RetryableTopic(attempts = "3", backoff = @Backoff(value = 3000L)) public void readMessages(String message) { throw new NullPointerException("消息内容为空"); } @DltHandler public void processDltMessages(ConsumerRecord<String, String> record) { String message = record.value(); // 从消息头中提取异常信息 String exceptionMessage = new String( record.headers().lastHeader(KafkaHeaders.EXCEPTION_MESSAGE).value(), StandardCharsets.UTF_8 ); String stackTrace = new String( record.headers().lastHeader(KafkaHeaders.EXCEPTION_STACKTRACE).value(), StandardCharsets.UTF_8 ); // 自定义逻辑处理 if (exceptionMessage.startsWith("NullPointerException")) { System.out.println("处理空内容的DLT消息:" + message); } }
补充说明
- 默认情况下,Spring Kafka会自动将原始异常的消息文本和栈跟踪信息分别存入
KafkaHeaders.EXCEPTION_MESSAGE和KafkaHeaders.EXCEPTION_STACKTRACE这两个消息头中,无需额外配置。 - 如果需要获取完整的
Throwable对象,你可以自定义RetryTopicConfiguration来配置异常序列化器,但这种场景较少,通常消息头中的信息已足够满足业务需求。
内容的提问来源于stack exchange,提问作者pvpkiran
相关产品推荐
相关产品推荐

