如何在死信队列DLQ监听器中识别RabbitMQ错误处理器抛出的异常类型?
解决方案
核心思路是:异常信息不会自动传递到死信队列的消息头中,你可以在自定义错误处理器中手动把异常类型/定制化重试参数写入原生AMQP消息的头部,消息进入死信队列时会携带这些自定义头,在死信队列处理器中读取对应头值即可动态调整重试策略。
第一步:修改MyErrorHandler,写入自定义消息头
在抛出可重试异常前,把异常类型、定制化延迟值写入消息头,消息进入死信队列时会保留这些头信息:
public class MyErrorHandler implements RabbitListenerErrorHandler { // 自定义消息头key,避免和RabbitMQ原生头冲突 public static final String X_EXCEPTION_TYPE = "x-exception-type"; public static final String X_CUSTOM_RETRY_DELAY = "x-custom-retry-delay"; @Override public Object handleError(Message amqpMessage, org.springframework.messaging.Message<?> message, ListenerExecutionFailedException exception) { Throwable cause = exception.getCause(); // 检查异常是致命还是可重试 if (cause instanceof FatalException || cause instanceof RetriesExceededException) { return new Status("FAIL!"); } // 写入异常类型到消息头 amqpMessage.getMessageProperties().getHeaders() .put(X_EXCEPTION_TYPE, cause.getClass().getSimpleName()); // 限流异常直接写入更大的重试延迟,也可以统一放到DLQ中判断 if (cause instanceof RateLimitException) { amqpMessage.getMessageProperties().getHeaders() .put(X_CUSTOM_RETRY_DELAY, properties.getRateLimitRetryDelay()); } else { amqpMessage.getMessageProperties().getHeaders() .put(X_CUSTOM_RETRY_DELAY, properties.getRetryDelay()); } // 抛出异常让消息NACK进入死信队列,此时消息已携带自定义头 throw exception; } }
第二步:修改MyDLQueueHandler,读取自定义头动态调整延迟
从死信消息的头中读取之前写入的异常标识/定制延迟值即可:
@Service public class MyDLQueueHandler { private final MyAppProperties properties; private final MessageHeaderProcessor headerProcessor; private final RabbitProducerService rabbitProducerService; public MyDLQueueHandler(MyProperties properties, MessageHeaderProcessor headerProcessor, RabbitProducerService rabbitProducerService) { this.properties = properties; this.headerProcessor = headerProcessor; this.rabbitProducerService = rabbitProducerService; } /** * 由于仲裁队列不支持消息TTL,手动监听死信队列并带延迟重发消息,实现消息重试逻辑 */ @RabbitListener(queues = {"#{myDLQueue.getName()}"}) public void handleError(@Payload Object message, @Headers MessageHeaders headers) { String routingKey = headerProcessor.getRoutingKey(headers); // 读取自定义延迟,无值则用默认配置兜底 Integer retryDelay = headers.get(MyErrorHandler.X_CUSTOM_RETRY_DELAY, Integer.class); if (retryDelay == null) { retryDelay = properties.getRetryDelay(); } // 也可以读取异常类型自行判断延迟,两种方式二选一: // String exceptionType = headers.get(MyErrorHandler.X_EXCEPTION_TYPE, String.class); // if (RateLimitException.class.getSimpleName().equals(exceptionType)) { // retryDelay = properties.getRateLimitRetryDelay(); // } Map<String, Object> newHeaders = Map.of( MessageHeaderProcessor.DELAY, retryDelay, MessageHeaderProcessor.RETRIES_HEADER, headerProcessor.getRetries(headers) + 1 ); rabbitProducerService.sendMessageDelayed(message, routingKey, newHeaders); } }
你也可以根据业务需要扩展自定义头的内容,比如限流接口返回的官方建议等待时间、异常栈摘要等,只要是可序列化的类型都可以放入AMQP消息头传递。
内容的提问来源于stack exchange,提问作者Süleyman Gezsat
相关产品推荐
相关产品推荐

