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

如何配置ActiveMQ Artemis实现AMQP与OpenWire跨协议消息互通

ActiveMQ Artemis跨协议主题消息互通问题

服务器配置

public static void updateConfig(Configuration config) {
    try {
        config.setPersistenceEnabled(false)
              .setSecurityEnabled(false)
              .addAcceptorConfiguration("tcp", "tcp://localhost:61616")
              .addAcceptorConfiguration("amqp", "tcp://localhost:5672");
    } catch (Exception ex) {
        ex.printStackTrace();
    }
}

问题现象

服务器同时支持AMQP(NodeJS客户端)和OpenWire(Java客户端)协议,使用主题TestTopic收发消息,但两类客户端消息无法互通:

  • AMQP发布的消息仅能被NodeJS订阅者接收,Java OpenWire订阅者收不到
  • Java客户端发布的消息仅能被Java订阅者接收,NodeJS订阅者收不到
    服务器后台可观测到两类协议的消息,但未向TestTopic的所有订阅者广播。

客户端配置

NodeJS订阅者配置

NodeJS客户端指定的目标主题为'topic://TestTopic':

// subscriber.js
var args = require('./options.js').options({
  'client': { default: 'my-client', describe: 'name of identifier for client container'},
  'subscription': { default: 'my-subscription', describe: 'name of identifier for subscription'},
  't': { alias: 'topic', default: 'topic://TestTopic', describe: 'name of topic to subscribe to'},
  'h': { alias: 'host', default: 'localhost', describe: 'dns or ip name of server where you want to connect'},
  'p': { alias: 'port', default: 5672, describe: 'port to connect to'}
}).help('help').argv;

var connection = require('rhea').connect({ port:args.port, host: args.host, container_id:args.client });
connection.on('receiver_open', function (context) {
  console.log('subscribed');
});
connection.on('message', function (context) {
  if (context.message.body === 'detach') {
      // detaching leaves the subscription active, so messages sent
      // while detached are kept until we attach again
      context.receiver.detach();
      context.connection.close();
  } else if (context.message.body === 'close') {
      // closing cancels the subscription
      context.receiver.close();
      context.connection.close();
  } else {
      console.log(context.message.body);
  }
});
// the identity of the subscriber is the combination of container id
// and link (i.e. receiver) name
connection.open_receiver({name:args.subscription, source:{address:args.topic, durable:2, expiry_policy:'never'}});

Java监听器配置

Java监听器指定的目标主题为TestTopic:

// client-sender.java
@JmsListener(destination = "TestTopic", selector = "${selector}")
public void receiveMessage(Message message) {
    String type = (message instanceof MapMessage) ? "MapMessage" : "TextMessage";
    try {
        logMessageBody(message);
        }
    } catch (Exception ex) {
        logger.error("Exception caught while logging the received message: ");
        logger.error(ex + ": " + ex.getCause());
    }
}

问题原因

官方文档明确说明:

应为目标地址添加queue://前缀以使用队列类型目标,或添加topic://前缀以使用主题类型目标。若省略前缀,默认类型为队列。

当前两类客户端的地址前缀不一致:

  • NodeJS客户端使用带topic://前缀的地址topic://TestTopic,对应服务器上的独立地址
  • Java JMS客户端使用无前缀的TestTopic,ActiveMQ Artemis对JMS主题会自动映射到地址TestTopic

两者实际订阅的是服务器上的不同地址,因此无法收到彼此的消息。

解决方案

任选以下一种方式统一地址即可:

方式1:Java客户端添加前缀

修改Java监听器的目标地址,添加topic://前缀:

@JmsListener(destination = "topic://TestTopic", selector = "${selector}")
public void receiveMessage(Message message) {
    // 原有业务逻辑保持不变
}

方式2:NodeJS客户端移除前缀

修改NodeJS客户端的主题配置,去掉topic://前缀:

't': { alias: 'topic', default: 'TestTopic', describe: 'name of topic to subscribe to'},

同时确保NodeJS发布消息时也使用无前缀的TestTopic地址,服务器会自动识别为主题类型。

验证

修改完成后,两类客户端将订阅服务器上的同一个地址,消息即可在所有订阅者之间正常互通。

内容的提问来源于stack exchange,提问作者Potion

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 10:15:56