You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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接收到的消息内容也与此一致。请问:

  1. 能否在smallrye发送消息时处理,使消息内容可读?
  2. 如何用activemq-client正确解析这类消息?
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 23:33:15