Quarkus基于SmallRye与Spring Boot ActiveMQ通信的消息格式问题
问题
我负责一个基于Quarkus的应用,需通过ActiveMQ与另一Spring Boot应用通信。Quarkus端使用smallrye库对接ActiveMQ,对方Spring Boot应用则使用org.apache.activemq:activemq-client库。我无法修改对方应用,但可协调做小幅度调整。
我编写示例代码尝试双向通信,Spring Boot端发送的消息可被Quarkus端正常接收,但反向发送时,对方收到的消息正文开头混入了元数据(可识别部分属性名),ActiveMQ管理界面中查看smallrye发送的消息也存在类似情况。
发送消息的代码如下:
@Outgoing("outgoing-queue") public Multi<Message<String>> sendWithSomeMetadataSet() { return Multi.createFrom() .ticks().every(Duration.ofSeconds(20)) .map(x -> Message.of("Hello")) .map(m -> m.addMetadata(createMetadata("Test"))); } static OutgoingAmqpMetadata createMetadata(String address) { return OutgoingAmqpMetadata.builder() .withAddress(address) .withApplicationProperty("foo", "bar") .withApplicationProperty("brown", "fox") .withDurable(true) .withCreationTime(Instant.now().toEpochMilli()) .withExpiryTime(Instant.now().plus(12, ChronoUnit.HOURS).toEpochMilli()) .build(); }
在ActiveMQ管理界面中看到的消息内容如下:
Sp�APSs�%@@�Test@@@@@���;bm����4m@CSt��foo�bar�brown�foxSw�Hello
Spring Boot端用activemq-client接收到的消息内容也与此一致。请问:
- 能否在smallrye发送消息时处理,使消息内容可读?
- 如何用activemq-client正确解析这类消息?
- Quarkus是否有可替代smallrye的库?
解决方案
1. 在SmallRye端调整发送逻辑,让消息正文可读
问题根源是SmallRye默认用AMQP协议序列化消息,而Spring Boot的activemq-client默认期望OpenWire协议的消息格式,两者不兼容导致元数据被混入正文。调整方案如下:
- 放弃使用
OutgoingAmqpMetadata,改用Quarkus提供的JMS兼容API发送消息,确保消息以OpenWire格式序列化。 - 替换后的代码示例:
import jakarta.enterprise.context.ApplicationScoped; import jakarta.jms.ConnectionFactory; import jakarta.jms.JMSContext; import io.smallrye.mutiny.Multi; import java.time.Duration; @ApplicationScoped public class MessageSender { private final ConnectionFactory connectionFactory; public MessageSender(ConnectionFactory connectionFactory) { this.connectionFactory = connectionFactory; } @Outgoing("outgoing-queue") public Multi<String> sendMessage() { return Multi.createFrom() .ticks().every(Duration.ofSeconds(20)) .map(x -> { try (JMSContext context = connectionFactory.createContext()) { // 构造JMS消息并设置原元数据对应的属性 var message = context.createTextMessage("Hello"); message.setStringProperty("foo", "bar"); message.setStringProperty("brown", "fox"); message.setJMSDeliveryMode(2); // 对应durable=true message.setJMSTimestamp(System.currentTimeMillis()); message.setJMSExpiration(System.currentTimeMillis() + 12 * 3600 * 1000); // 发送到目标队列 context.createProducer().send(context.createQueue("Test"), message); return "Hello"; } }); } }
- 确保Quarkus依赖中引入
quarkus-jms-activemq扩展,替代原AMQP相关的SmallRye扩展,自动适配OpenWire协议。
2. Spring Boot端用activemq-client解析现有消息(临时方案)
若无法调整Quarkus端,可在Spring Boot端切换为AMQP客户端解析消息:
- 引入
org.apache.qpid:qpid-jms-client依赖,用AMQP协议接收消息,提取正文和属性:
import org.apache.qpid.jms.JmsConnectionFactory; import jakarta.jms.Message; import jakarta.jms.TextMessage; public class AmqpMessageReceiver { public void receiveMessage() throws Exception { JmsConnectionFactory factory = new JmsConnectionFactory("tcp://activemq-host:5672"); try (jakarta.jms.Connection connection = factory.createConnection()) { connection.start(); jakarta.jms.Session session = connection.createSession(false, jakarta.jms.Session.AUTO_ACKNOWLEDGE); jakarta.jms.MessageConsumer consumer = session.createConsumer(session.createQueue("Test")); Message message = consumer.receive(); if (message instanceof TextMessage textMessage) { String body = textMessage.getText(); String foo = textMessage.getStringProperty("foo"); // 处理正文和属性 } } } }
- 此方案需修改Spring Boot端的依赖和接收逻辑,仅适合对方能做小幅度调整的场景,优先推荐Quarkus端改JMS方案。
3. Quarkus的替代库
Quarkus支持直接使用标准JMS API对接ActiveMQ,替代SmallRye AMQP的方案:
- 使用
quarkus-jms-activemq扩展:基于OpenWire协议,与Spring Boot的activemq-client完全兼容,彻底避免协议格式冲突。 - 引入依赖(pom.xml):
<dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-jms-activemq</artifactId> </dependency>
- 该扩展会自动配置ActiveMQ连接工厂,直接使用JMS API发送/接收消息,无需额外处理协议适配。
内容的提问来源于stack exchange,提问作者somebody
相关产品推荐
相关产品推荐

