如何在Spring Cloud Stream(RabbitMQ)中合理封装含DLQ的重试逻辑?
消息投递次数检查的公共组件实现疑问
官方推荐的异常重入队处理方式
根据官方文档,异常场景下消息重入队的推荐实现代码如下:
@Bean public Consumer<Message<String>> listen() { return message -> { Map<?,?> death = message.getHeaders().get("x-death"); if (death != null && death.get("count").equals(3L)) { // 放弃处理,不发送到死信交换机 throw new ImmediateAcknowledgeAmqpException("Failed after 4 attempts"); } throw new AmqpRejectAndDontRequeueException("failed"); }; }
我可以修改消费者代码抛出这些异常,但当前负责的代码库包含多个微服务,这种检查投递次数的逻辑会产生大量重复代码。想咨询是否能将这部分逻辑提取为可放入公共库的独立组件?
自行实现的方案
因未得到回复,我自己实现了以下方案:
@GlobalChannelInterceptor public class DeliveryAttemptsCheckingInterceptor implements org.springframework.messaging.support.ChannelInterceptor { @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { if (message.getHeaders().containsKey("x-death") && xDeathCount(message) >= MAX_DELIVERY_ATTEMPTS) { throw new ImmediateAcknowledgeAmqpException(""); // 直接丢弃消息 } else { return message; } } }
该方案可以正常工作,但我不确定这种使用ChannelInterceptor的方式是否属于滥用。
内容的提问来源于stack exchange,提问作者FoxMulder
相关产品推荐
相关产品推荐

