Spring 的 RabbitMQ 直接回复问题
Problem with RabbitMQ Direct reply-to with Spring
我正在开发一个将消息发送到服务器的应用程序,然后修改给定的消息并使用 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 协议上的 SEND 帧作为 JSON 发送并转换。然后一个新队列
动态创建并添加到 MessageListenerContainer。因此,当消息到达代理时,我想在服务器端对其进行修改并发送回 amq.rabbitmq.reply-to
并将原始消息发送到在 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]
这是我在 PracticalMessageListener
中调用 rabbitTemplate.convertAndSend
时的地方
编辑
我摆脱了 DirectMessageListenerContainer
中的设置 amq.rabbitmq.reply-to
并实现了 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()
是在运行时定义的路由键,通过 selecting 接收器,例如如果我 select '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
放在单独的 class 中并将其 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;
}
}
是的,您误解了该功能。
每个通道都有自己的伪队列;你只能从同一个频道接收,所以一般的消息监听器容器不会破解它。
directMessageListenerContainer.setQueueNames("amq.rabbitmq.reply-to");
你根本做不到。
该框架已经在 RabbitTemplate
内部直接支持直接回复。 RabbitTemplate
有自己的 DirectReplyToMessageListenerContainer
维护频道池。
每个请求都会检查一个通道并在那里返回回复,然后通道返回到池中以供另一个请求重用。
使用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
我正在开发一个将消息发送到服务器的应用程序,然后修改给定的消息并使用 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 协议上的 SEND 帧作为 JSON 发送并转换。然后一个新队列
动态创建并添加到 MessageListenerContainer。因此,当消息到达代理时,我想在服务器端对其进行修改并发送回 amq.rabbitmq.reply-to
并将原始消息发送到在 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]
这是我在 PracticalMessageListener
rabbitTemplate.convertAndSend
时的地方
编辑
我摆脱了 DirectMessageListenerContainer
中的设置 amq.rabbitmq.reply-to
并实现了 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()
是在运行时定义的路由键,通过 selecting 接收器,例如如果我 select '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
放在单独的 class 中并将其 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;
}
}
是的,您误解了该功能。
每个通道都有自己的伪队列;你只能从同一个频道接收,所以一般的消息监听器容器不会破解它。
directMessageListenerContainer.setQueueNames("amq.rabbitmq.reply-to");
你根本做不到。
该框架已经在 RabbitTemplate
内部直接支持直接回复。 RabbitTemplate
有自己的 DirectReplyToMessageListenerContainer
维护频道池。
每个请求都会检查一个通道并在那里返回回复,然后通道返回到池中以供另一个请求重用。
使用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