Spring Boot整合ActiveMQ Artemis收发消息时解码异常求助
我有一个基于Spring Boot 2.7.6的ActiveMQ Artemis消费者应用,代码如下:
package broker.consumer; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.apache.activemq.broker.BrokerService; import org.springframework.jms.annotation.JmsListener; @SpringBootApplication public class Application { public static void main(String[] args) { SpringApplication.run(Application.class, args); } @Bean public BrokerService broker() throws Exception { BrokerService broker = new BrokerService(); broker.addConnector("tcp://localhost:61616"); broker.setPersistent(false); broker.start(); return broker; } @JmsListener(destination = "foo") public void listen(String in) { System.out.println(in); } }
另有一个Spring Boot生产者应用,用于向Broker地址发送消息,代码如下:
package broker.producer; import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.jms.core.JmsTemplate; import org.springframework.jms.support.destination.JndiDestinationResolver; import org.springframework.stereotype.Service; @Service public class JmsProducer { @Value("${spring.artemis.broker-url}") private String brokerUrl; @Value("${spring.jms.template.default-destination}") private String defaultDestination; Logger log = LoggerFactory.getLogger(JmsProducer.class); @Bean public ActiveMQConnectionFactory activeMQConnectionFactory() { log.info("BrokerUrl: {}", brokerUrl); ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(brokerUrl); return activeMQConnectionFactory; } @Bean public JndiDestinationResolver jndiDestinationResolver() { return new JndiDestinationResolver(); } @Bean public JmsTemplate jmsTemplate() { JmsTemplate template = new JmsTemplate(); template.setConnectionFactory(activeMQConnectionFactory()); template.setPubSubDomain(false); // false for a Queue, true for a Topic template.setDefaultDestinationName(defaultDestination); return template; } public void send(String message) { JmsTemplate jmsTemplate = jmsTemplate(); log.info("Sending message='{}'", message); jmsTemplate.convertAndSend(message); log.info("Sent message='{}'", message); } }
两个应用使用相同的application.properties配置:
spring.artemis.mode=EMBEDDED spring.artemis.broker-url=tcp://localhost:61616 spring.artemis.user=admin spring.artemis.password=secret spring.artemis.embedded.enabled=true spring.jms.template.default-destination=my-queue-1
启动两个应用并调用send方法发送消息时,生产者应用出现如下错误:
2024-01-15 15:50:28.462 ERROR 1012 --- [-netty-threads)] org.apache.activemq.artemis.core.client : AMQ214013: Failed to decode packet
java.lang.IllegalArgumentException: AMQ219032: Invalid type: 1
at org.apache.activemq.artemis.core.protocol.core.impl.PacketDecoder.decode(PacketDecoder.java:499) ~[artemis-core-client-2.19.1.jar:2.19.1]
消费者应用则出现如下警告:
2024-01-15 15:50:28.464 WARN 986 --- [0.1:57714@61616] o.a.a.b.TransportConnection.Transport : Transport Connection to: tcp://127.0.0.1:57714 failed: Unknown data type: 77
问题核心是协议不兼容:
- 消费者嵌入的是ActiveMQ Classic的
BrokerService(org.apache.activemq.broker包),默认用OpenWire协议 - 生产者用的是ActiveMQ Artemis客户端(
org.apache.activemq.artemis.jms.client包),默认用Artemis原生Core协议
两者协议不匹配导致数据包解码失败,同时消费者监听的目标"foo"和生产者发送的"my-queue-1"不一致,也会导致消息无法接收。以下两种方法二选一即可:
方法一:将消费者Broker替换为ActiveMQ Artemis
把Classic的BrokerService换成Artemis的嵌入式Broker,同时修正监听目标:
package broker.consumer; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ; import org.springframework.jms.annotation.JmsListener; @SpringBootApplication public class Application { public static void main(String[] args) { SpringApplication.run(Application.class, args); } @Bean public EmbeddedActiveMQ embeddedActiveMQ() throws Exception { EmbeddedActiveMQ broker = new EmbeddedActiveMQ(); broker.addConnector("tcp://localhost:61616"); broker.setPersistent(false); broker.start(); return broker; } @JmsListener(destination = "my-queue-1") // 修正为与生产者一致的目标 public void listen(String in) { System.out.println(in); } }
方法二:让Artemis生产者使用OpenWire协议
修改生产者的连接工厂,强制使用OpenWire协议,同时修正消费者监听目标:
// 生产者ActiveMQConnectionFactory配置修改 @Bean public ActiveMQConnectionFactory activeMQConnectionFactory() { log.info("BrokerUrl: {}", brokerUrl); ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(brokerUrl); activeMQConnectionFactory.setProtocolType("OPENWIRE"); // 强制指定OpenWire协议 return activeMQConnectionFactory; }
同时将消费者的@JmsListener(destination = "foo")改为@JmsListener(destination = "my-queue-1")。
内容的提问来源于stack exchange,提问作者Franek

