【问题标题】:How to work with Dead letter queues with reactor rabbitmq如何使用 reactor rabbitmq 处理死信队列
【发布时间】:2020-11-30 05:47:14
【问题描述】:

在我们的应用程序中,我们正在从 spring-amqp 迁移到 reactor-rabbitmq,以更好地适应应用程序的反应性质。我们一直在阅读来自 project reactor 的official guide。但是我不太确定,一旦重试用尽,如何将消息发送到死信队列。

在之前的实现中,我们抛出 Spring-AMPQ 提供的AmqpRejectAndDontRequeueException,这会导致消息自动进入死信队列,而不使用死信队列的专用发布者。如何在 reactor-rabbitmq 中做类似的事情?我是否需要编写一个专用的发布者,在重试用尽后从侦听器中调用它,或者有其他方法来处理它。此外,是否有关于 DLQ 和停车队列的官方项目反应器文档。

以下是两个版本的一些代码示例:

AMQP 版本:

@AllArgsConstructor
public class SampleListener {
    private static final Logger logger = LoggerFactory.getLogger(SampleListener.class);
    private final MessageListenerContainerFactory messageListenerContainerFactory;
    private final Jackson2JsonMessageConverter converter;

    @PostConstruct
    public void subscribe() {
        var mlc = messageListenerContainerFactory
                .createMessageListenerContainer(SAMPLE_QUEUE);
        MessageListener messageListener = message -> {
            try {
                TraceableMessage traceableMessage = (TraceableMessage) converter.fromMessage(message);
                ObjectMapper mapper = new ObjectMapper();
                mapper.registerModule(new JavaTimeModule());
                MyModel myModel = mapper.convertValue(traceableMessage.getMessage(), MyModel.class);
                MDC.put(CORRELATION_ID, traceableMessage.getCorrelationId());
                logger.info("Received message for id : {}", myModel.getId());
                processMessage(myModel)
                        .subscriberContext(ctx -> {
                            Optional<String> correlationId = Optional.ofNullable(MDC.get(CORRELATION_ID));
                            return correlationId.map(id -> ctx.put(CORRELATION_ID, id))
                                    .orElseGet(() -> ctx.put(CORRELATION_ID, UUID.randomUUID().toString()));
                        }).block();
                MDC.clear();
            } catch (Exception e) {
                logger.error(e.getMessage(), e);
                throw new AmqpRejectAndDontRequeueException(e.getMessage(), e);
            }
        };
        mlc.setupMessageListener(messageListener);
        mlc.start();
    }

processMessage 正在执行业务逻辑,如果失败,我想将其移至 DLQ。在 AMQP 的情况下工作正常。

Reactor RabbitMQ 版本:

@AllArgsConstructor
public class SampleListener {
    private static final Logger logger = LoggerFactory.getLogger(SampleListener.class);
    private final MessageListenerContainerFactory messageListenerContainerFactory;
    private final Jackson2JsonMessageConverter converter;

    @PostConstruct
    public void subscribe() {
        receiver.consumeAutoAck(SAMPLE_QUEUE)
                .subscribe(delivery -> {
                            TraceableMessage traceableMessage = Serializer.to(delivery.getBody(), TraceableMessage.class);
                            Mono.just(traceableMessage)
                                    .map(this::extractMyModel)
                                    .doOnNext(myModel -> logger.info("Received message for id : {}", myModel.getId()))
                                    .flatMap(this::processMessage)
                                    .doFinally(signalType -> MDC.clear())
                                    .retryWhen(Retry
                                            .fixedDelay(1, Duration.ofMillis(10000))
                                            .onRetryExhaustedThrow() //Move to DLQ
                                            .doAfterRetry(retrySignal -> {
                                                if ((retrySignal.totalRetries() + 1) >= 1) {
                                                    logger.info("Exhausted retries");
                                                    //Move to DLQ
                                                }
                                            }))
                                    .subscriberContext(ctx -> ctx.put(CORRELATION_ID, traceableMessage.getCorrelationId()))
                                    .subscribe();
                        }
                );

    }

comment//Move to DLQ 的这两个地方之一是我猜测消息应该发送到 DLQ 的哪个位置。这就是我决定我不能再处理这个的地方。如果有不同的发布者推送到 DLQ 或任何特定设置可以自动处理它。

请告诉我。

【问题讨论】:

  • 给我们看一些代码??
  • @PrashantPandey 添加了两者的代码版本。不确定它会有多大帮助,这就是我之前没有添加它的原因。

标签: java spring rabbitmq project-reactor reactor-rabbitmq


【解决方案1】:

我得到了这个问题的答案。我正在使用异步侦听器并使用consumeAutoAck。当我切换到consumeManualAck 时,我会得到一个AcknowledgebleDelivery,我可以在其中执行一个nack(false),这应该将其移至死信队列。

【讨论】:

  • 嗨!在这种情况下,您是如何使用 reactor rabbit 设置重试策略的?我有一个消费者,在异常情况下我会做nack(false),并且消息会立即发送到 DLQ。如何定义我希望消息在消费者中尝试处理多少次,直到它实际发送到 DLQ?
  • 嗨@hideburn,您可以查看我们的实现here。
  • 嗨@Mritunjay,谢谢,非常感谢!
猜你喜欢
  • 1970-01-01
  • 2013-11-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-08-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多