如何配置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
相关产品推荐
相关产品推荐

