java - RabbitMQ 使用 Spring 直接回复的问题

标签 java spring rabbitmq spring-rabbit

我正在开发一个将消息发送到服务器的应用程序,然后修改给定的消息并使用 Direct Reply-to 发送回 amq.rabbitmq.reply-to 队列。我已经遵循了教程 https://www.rabbitmq.com/direct-reply-to.html 但在实现它时遇到了一些问题。就我而言,据我了解,我需要以无确认模式使用来自伪队列 amq.rabbitmq.reply-to 的消息,在我的情况下是 MessageListenerContainer 。这是我的配置:

@Bean
    public Jackson2JsonMessageConverter messageConverter() {    
        ObjectMapper mapper = new ObjectMapper();
        return new Jackson2JsonMessageConverter(mapper);
    }

    @Bean
    public RabbitAdmin rabbitAdmin(ConnectionFactory connectionFactory) {
        return new RabbitAdmin(connectionFactory);
    }

    @Bean
    public RabbitTemplate rabbitTemplate(final ConnectionFactory connectionFactory) {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
        rabbitTemplate.setMessageConverter(messageConverter());
        rabbitTemplate.setReplyAddress("amq.rabbitmq.reply-to");
        return rabbitTemplate;
    }

    @Bean
    MessageListenerContainer messageListenerContainer(ConnectionFactory connectionFactory ) {       
        DirectMessageListenerContainer directMessageListenerContainer = new DirectMessageListenerContainer();
        directMessageListenerContainer.setConnectionFactory(connectionFactory);
        directMessageListenerContainer.setAcknowledgeMode(AcknowledgeMode.NONE);
        directMessageListenerContainer.setQueueNames("amq.rabbitmq.reply-to");
        directMessageListenerContainer.setMessageListener(new PracticalMessageListener());
        return directMessageListenerContainer;

    }

消息通过 STOM 协议(protocol)上的 SEND 帧以 JSON 形式发送并进行转换。然后是一个新队列
动态创建并添加到 MessageListenerContainer。因此,当消息到达代理时,我想在服务器端修改它并发送回 amq.rabbitmq.reply-to 并将原始消息发送到在 STOMP 中的 SUBSCRIBE 帧上订阅的路由键 messageTemp.getTo()

  @MessageMapping("/private")
  public void send2(MessageTemplate messageTemp) throws Exception {
      MessageTemplate privateMessage = new MessageTemplate(messageTemp.getPerson(),
              messageTemp.getMessage(), 
              messageTemp.getTo());

     AbstractMessageListenerContainer abstractMessageListenerContainer =
              (AbstractMessageListenerContainer) mlc;

       // here's the queue added to listener container   
      abstractMessageListenerContainer.addQueueNames(messageTemp.getTo());

      MessageProperties mp = new MessageProperties();
      mp.setReplyTo("amq.rabbitmq.reply-to");
      mp.setCorrelationId("someId");

      Jackson2JsonMessageConverter smc = new Jackson2JsonMessageConverter();
      Message message = smc.toMessage(messageTemp, mp);


      rabbitTemplate.sendAndReceive( 
              messageTemp.getTo() , message);
  }

当消息发送到messageTemp.getTo()路由键时,消息被修改onMessage方法

@Component
public class PracticalMessageListener implements MessageListener {

    @Autowired
    RabbitTemplate rabbitTemplate;

    @Override
    public void onMessage(Message message) {
        System.out.println(("message listener.."));
        String body = "{ \"processing\": \"123456789\"}";
       MessageProperties properties = new MessageProperties();

       // some business logic on the message body

        properties.setCorrelationId(message.getMessageProperties().getCorrelationId());
        Message responseMessage = new Message(body.getBytes(), properties);

        rabbitTemplate.convertAndSend("", 
                message.getMessageProperties().getReplyTo(), responseMessage);
    }

我可能误解了直接回复的概念以及文档:

Consume from the pseudo-queue amq.rabbitmq.reply-to in no-ack mode. There is no need to declare this "queue" first, although the client can do so if it wants.

问题是我需要从该队列的哪里消费?如果出现错误,我如何访问修改后的消息:

2020-01-15 22:17:09.688  WARN 96222 --- [pool-1-thread-5] s.a.r.l.ConditionalRejectingErrorHandler : Execution of Rabbit message listener failed.

org.springframework.amqp.rabbit.support.ListenerExecutionFailedException: Listener threw exception
Caused by: java.lang.NullPointerException: null
    at com.patrykmaryn.spring.second.PracticalMessageListener.onMessage(PracticalMessageListener.java:50) ~[classes/:na]

这是来 self 在rabbitTemplate.convertAndSend中调用PracticalMessageListener时的地方

编辑

我摆脱了在 amq.rabbitmq.reply-to 中设置 DirectMessageListenerContainer 并实现了 DirectReplyToMessageListenerContainer :

@Bean
    DirectReplyToMessageListenerContainer drtmlc (ConnectionFactory connectionFactory) {
        DirectReplyToMessageListenerContainer drtmlc =
                new DirectReplyToMessageListenerContainer(connectionFactory);
        drtmlc.setConnectionFactory(connectionFactory);
        drtmlc.setAcknowledgeMode(AcknowledgeMode.NONE);
        drtmlc.setMessageListener(new DirectMessageListener());
        return drtmlc;
    }

问题一定出在 onMessage 方法中,该方法不允许调用 rabbitTemplate 上的任何发送方法,我已经尝试过不同的现有路由键和交换。监听来自使用路由键 messageTemp.getTo() 定义的队列。

@Override
    public void onMessage(Message message) {
        System.out.println(("message listener.."));

        String receivedRoutingKey = message.getMessageProperties()
           .getReceivedRoutingKey();
        System.out.println(" This is received routingkey: " + 
            receivedRoutingKey);

           /// ..... rest of code goes here

        rabbitTemplate.convertAndSend("", 
                message.getMessageProperties().getReplyTo(), 
                responseMessage);

其中 messageTemp.getTo() 是在运行时定义的路由键,通过选择接收器,例如,如果我选择“user1”,它将打印出“user1”。

这是第一次尝试发送消息:

2020-01-16 02:22:20.213 DEBUG 28490 --- [nboundChannel-6] .WebSocketAnnotationMethodMessageHandler : Searching methods to handle SEND /app/private session=45yca5sy, lookupDestination='/private'
2020-01-16 02:22:20.214 DEBUG 28490 --- [nboundChannel-6] .WebSocketAnnotationMethodMessageHandler : Invoking PracticalTipSender#send2[1 args]
2020-01-16 02:22:20.239  INFO 28490 --- [nboundChannel-6] o.s.a.r.l.DirectMessageListenerContainer : SimpleConsumer [queue=user1, consumerTag=amq.ctag-Evyiweew4C-K1mXmy2XqUQ identity=57b19488] started
2020-01-16 02:22:20.268  INFO 28490 --- [nboundChannel-6] o.s.s.c.ThreadPoolTaskScheduler          : Initializing ExecutorService
2020-01-16 02:22:20.269  INFO 28490 --- [nboundChannel-6] .l.DirectReplyToMessageListenerContainer : Container initialized for queues: [amq.rabbitmq.reply-to]
2020-01-16 02:22:20.286  INFO 28490 --- [nboundChannel-6] .l.DirectReplyToMessageListenerContainer : SimpleConsumer [queue=amq.rabbitmq.reply-to, consumerTag=amq.ctag-IXWf-zEyI34xzQSSfbijzg identity=4bedbba5] started

第二个失败:

2020-01-16 02:23:20.247 DEBUG 28490 --- [nboundChannel-3] .WebSocketAnnotationMethodMessageHandler : Searching methods to handle SEND /app/private session=45yca5sy, lookupDestination='/private'
2020-01-16 02:23:20.248 DEBUG 28490 --- [nboundChannel-3] .WebSocketAnnotationMethodMessageHandler : Invoking PracticalTipSender#send2[1 args]
2020-01-16 02:23:20.248  WARN 28490 --- [nboundChannel-3] o.s.a.r.l.DirectMessageListenerContainer : Queue user1 is already configured for this container: org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer@3b152928, ignoring add
2020-01-16 02:23:20.250  WARN 28490 --- [nboundChannel-3] o.s.a.r.l.DirectMessageListenerContainer : Queue user1 is already configured for this container: org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer@3b152928, ignoring add
message listener..
 This is received routingkey: user1
2020-01-16 02:23:20.271  WARN 28490 --- [pool-1-thread-5] s.a.r.l.ConditionalRejectingErrorHandler : Execution of Rabbit message listener failed.

org.springframework.amqp.rabbit.support.ListenerExecutionFailedException: Listener threw exception

编辑

DirectReplyToMessageListenerContainer 放在单独的类中并将其 MessageListener 设置为 @Bean 并且 directMessageListenerContainer.setMessageListener(practicalMessageListener()); 作为 @Bean 似乎摆脱了 NPE。但即使回复到 amq.rabbitmq.reply-to.g2dkABVyYWJ..... ,它似乎也没有在 DirectReplyToMessageListenerContainer drtmlc 中被监听。

@Component
class DirectMessageListener implements MessageListener {
    // This doesn't get invoked...
    @Override
    public void onMessage(Message message) {
        System.out.println("direct reply message sent..");

    }
}

@Component
class ReplyListener {

    @Bean
    public DirectMessageListener directMessageListener() {
        return new DirectMessageListener(); 
    }

    @Bean
    DirectReplyToMessageListenerContainer drtmlc (ConnectionFactory connectionFactory) {
        DirectReplyToMessageListenerContainer drtmlc =
                new DirectReplyToMessageListenerContainer(connectionFactory);
        drtmlc.setConnectionFactory(connectionFactory);
        drtmlc.setAcknowledgeMode(AcknowledgeMode.NONE);
        drtmlc.setMessageListener(directMessageListener());
        return drtmlc;
    }
}

最佳答案

是的,您误解了该功能。

每个 channel 都有自己的伪队列;您只能从同一 channel 接收信息,因此通用消息监听器容器不会破解它。

directMessageListenerContainer.setQueueNames("amq.rabbitmq.reply-to");

你根本做不到。

该框架已经在 RabbitTemplate 内部支持直接回复。 RabbitTemplate 有自己的 DirectReplyToMessageListenerContainer,它维护着一个 channel 池。

每个请求都会检查一个 channel ,并在此处返回回复,然后该 channel 将返回到池中以供另一个请求重用。

使用RabbitTemplate.convertSendAndReceive();默认行为(在最近的版本中)将自动使用直接回复。

编辑

为什么不让框架完成所有繁重的工作,而您只需专注于业务逻辑:

@SpringBootApplication
public class So59760805Application {

    public static void main(String[] args) {
        SpringApplication.run(So59760805Application.class, args);
    }

    @Bean
    public SimpleMessageListenerContainer container(ConnectionFactory cf) {
        SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(cf);
        container.setQueueNames("foo");
        container.setMessageListener(new MessageListenerAdapter(new MyListener()));
        return container;
    }

    @Bean
    public MyExtendedTemplate template(ConnectionFactory cf) {
        return new MyExtendedTemplate(cf);
    }

    @Bean
    public ApplicationRunner runner(RabbitTemplate template) {
        return args -> System.out.println(template.convertSendAndReceive("", "foo", "test"));
    }

}

class MyListener {

    public String handleMessage(String in) {
        return in.toUpperCase();
    }

}

class MyExtendedTemplate extends RabbitTemplate {

    MyExtendedTemplate(ConnectionFactory cf) {
        super(cf);
    }

    @Override
    public void onMessage(Message message) {
        System.out.println("Response received (before conversion): " + message);
        super.onMessage(message);
    }

}

兔子模板默认使用直接回复(内部)。

Response received (before conversion): (Body:'TEST' MessageProperties [headers={}, correlationId=1, ...receivedRoutingKey=amq.rabbitmq.reply-to.g2dkAA5yYWJiaXRAZ29sbHVtMgAAeE0AAADmAw==.RQ/uxjR79PX/hZF+7iAdWw==, ...
TEST

关于java - RabbitMQ 使用 Spring 直接回复的问题,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/59760805/

相关文章:

spring - 使用 RabbitMQ 的 SimpMessagingTemplate.convertAndSend 工作速度很慢

rabbitmq - 如何启动 RabbitMQ 节点?

java - ReadExternalStorage 权限被拒绝,即使在 list 中声明它之后也是如此

java - Path.startsWith() 的奇怪结果

java - 如何在 java 中将 UUID 保存为二进制 (16)

java - Spring mvc 映射 header 和参数

RabbitMq:永久添加用户?

java - 无法在 Java 中解析此多部分 MIME 消息正文

spring - grails插件动态bean创建

java - 从 Spring Data 存储库获取支持的类型