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

Spring Boot整合ActiveMQ Artemis收发消息时解码异常求助

问题:ActiveMQ Artemis生产者与Classic Broker通信报错

我有一个基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 02:25:57