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

ActiveMQ Artemis多播转单播Divert失效及控制台预览问题求助

ActiveMQ Artemis 转发与消息显示问题

我尝试将multicast地址通过Divert转发至anycast地址,并通过Bridge将anycast数据转发至其他服务器,参考了官方divert示例。配置完成后遇到两个问题:

  • 使用JmsTemplate设置PubSubDomain为true向multicast发送消息时,multicast下的Queue可正常接收,但Divert未生效,桥接服务器无法收到消息;
  • 改用artemisClient发送消息时,multicast队列与桥接服务器均可接收消息,但web-console无法预览消息文本,提示“Unsupported message body type which cannot be displayed by hawtio”。

核心配置文件 broker.xml

<core xmlns="urn:activemq:core" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
      xsi:schemaLocation="urn:activemq:core">

    <name>subscribe-embedded-artemis</name>

    <persistence-enabled>true</persistence-enabled>

    <journal-type>MAPPED</journal-type>

    <paging-directory>prinect-data-flow/activemq-artemis/data/paging</paging-directory>

    <bindings-directory>prinect-data-flow/activemq-artemis/data/bindings</bindings-directory>

    <journal-directory>prinect-data-flow/activemq-artemis/data/journal</journal-directory>

    <large-messages-directory>prinect-data-flow/activemq-artemis/data/large-messages</large-messages-directory>

    <journal-datasync>true</journal-datasync>

    <journal-min-files>2</journal-min-files>

    <journal-pool-files>10</journal-pool-files>

    <journal-device-block-size>4096</journal-device-block-size>

    <journal-file-size>10M</journal-file-size>

    <journal-buffer-timeout>40000</journal-buffer-timeout>

    <journal-max-io>4096</journal-max-io>

    <!-- 磁盘使用量扫描间隔(毫秒) -->
    <disk-scan-period>5000</disk-scan-period>

    <!-- 磁盘使用率上限,达到后系统会阻塞或关闭不支持流控的连接 -->
    <max-disk-usage>90</max-disk-usage>

    <!-- 启用死锁等问题检测 -->
    <critical-analyzer>true</critical-analyzer>

    <critical-analyzer-timeout>120000</critical-analyzer-timeout>

    <critical-analyzer-check-period>60000</critical-analyzer-check-period>

    <critical-analyzer-policy>HALT</critical-analyzer-policy>


    <page-sync-timeout>1352000</page-sync-timeout>

    <connectors>
        <connector name="remote-connector">tcp://localhost:61617</connector>
    </connectors>

    <acceptors>
        <acceptor name="netty">tcp://0.0.0.0:61616?protocols=AMQP,CORE</acceptor>
    </acceptors>

    <diverts>
        <divert name="multicastToAnycastDivert">
            <routing-name>toAnycast</routing-name>
            <address>dataSubTopic</address>
            <forwarding-address>forwardMessageToCloud</forwarding-address>
            <exclusive>false</exclusive>
        </divert>
    </diverts>

    <bridges>
        <bridge name="my-bridge">
            <queue-name>forwardMessageToCloud</queue-name>
            <forwarding-address>remotePrinectDataForward</forwarding-address>
            <reconnect-attempts>-1</reconnect-attempts>
            <static-connectors>
                <connector-ref>remote-connector</connector-ref>
            </static-connectors>
        </bridge>
    </bridges>

    <security-settings>
        <security-setting match="#">
            <permission type="createNonDurableQueue" roles="guest"/>
            <permission type="deleteNonDurableQueue" roles="guest"/>
            <permission type="createDurableQueue" roles="guest"/>
            <permission type="deleteDurableQueue" roles="guest"/>
            <permission type="createAddress" roles="guest"/>
            <permission type="deleteAddress" roles="guest"/>
            <permission type="consume" roles="guest"/>
            <permission type="browse" roles="guest"/>
            <permission type="send" roles="guest"/>
            <!-- 允许manage权限,否则artemis数据导入工具无法工作 -->
            <permission type="manage" roles="guest"/>
        </security-setting>
    </security-settings>

    <address-settings>
        <address-setting match="#">
            <!-- 重发延迟乘数 -->
            <redelivery-delay-multiplier>1.5</redelivery-delay-multiplier>
            <!-- 初始重发延迟 -->
            <redelivery-delay>5000</redelivery-delay>
            <!-- 重发冲突避免因子 -->
            <redelivery-collision-avoidance-factor>0.15</redelivery-collision-avoidance-factor>
            <!-- 最大重发延迟 -->
            <max-redelivery-delay>50000</max-redelivery-delay>
            <dead-letter-address>DLA</dead-letter-address>
            <max-delivery-attempts>3</max-delivery-attempts>
            <auto-create-dead-letter-resources>true</auto-create-dead-letter-resources>
            <dead-letter-queue-prefix/>
            <dead-letter-queue-suffix>.DLQ</dead-letter-queue-suffix>

            <expiry-address>expiryAddress</expiry-address>
            <auto-create-expiry-resources>true</auto-create-expiry-resources>
            <expiry-queue-prefix/>
            <expiry-queue-suffix>.EXP</expiry-queue-suffix>
        </address-setting>
    </address-settings>

    <addresses>
        <address name="DLA">
            <anycast>
                <queue name="DLA"/>
            </anycast>
        </address>
        <address name="dataSubTopic">
            <multicast>
                <queue name="dataSubQueue1"/>
                <queue name="dataSubQueue2"/>
            </multicast>
        </address>
        <address name="forwardMessageToCloud">
            <anycast>
                <queue name="forwardMessageToCloud"/>
            </anycast>
        </address>
    </addresses>
</core>

发送代码

public void jmsSend(String message) {
    jmsTemplate.setPubSubDomain(true);
    jmsTemplate.convertAndSend("dataSubTopic", message);
}

public void clientSend(String message) {
    try (ClientSession session = artemisServerLocatorConfig.serverLocator()
            .createSessionFactory()
            .createSession()) {

        session.start();
        ClientProducer producer = session.createProducer("dataSubTopic");
        ClientMessage clientMessage = session.createMessage(true);
        clientMessage.setType(Message.TEXT_TYPE);
        clientMessage.getBodyBuffer()
                .writeString(message);
        producer.send(clientMessage);
    } catch (Exception e) {
        log.error("message = {}", e.getMessage());
    }
}

桥接服务器消息显示问题截图

Bridge Server Message


解决方案

1. JmsTemplate发送时Divert不生效的问题

原因:JmsTemplate使用JMS API发送消息时,PubSubDomain=true会将消息发送到Artemis自动生成的JMS Topic地址(默认前缀为jms.topic.),而当前Divert配置的address是原生Core地址dataSubTopic,无法匹配JMS消息的实际地址。

解决方法:

  • 修改Divert配置,匹配JMS Topic的实际地址:
<divert name="multicastToAnycastDivert">
    <routing-name>toAnycast</routing-name>
    <address>jms.topic.dataSubTopic</address>
    <forwarding-address>forwardMessageToCloud</forwarding-address>
    <exclusive>false</exclusive>
</divert>
  • 或者配置JmsTemplate的ConnectionFactory,禁用JMS地址前缀,让消息直接发送到dataSubTopic:
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
connectionFactory.setUseJmsMessageIDAsCoreMessageID(true);
connectionFactory.setPrefixJmsDestinations(false);
jmsTemplate.setConnectionFactory(connectionFactory);

2. artemisClient发送时web-console无法预览消息的问题

原因:直接通过Core Client写入字符串到消息体时,没有按照Hawtio可识别的格式封装,Hawtio需要消息符合JMS TextMessage的结构规范才能解析显示。

解决方法:

  • 调整Core Client代码,构造符合JMS TextMessage格式的消息:
import java.nio.charset.StandardCharsets;

public void clientSend(String message) {
    try (ClientSession session = artemisServerLocatorConfig.serverLocator()
            .createSessionFactory()
            .createSession()) {

        session.start();
        ClientProducer producer = session.createProducer("dataSubTopic");
        // 构造可被Hawtio识别的文本消息
        ClientMessage clientMessage = session.createMessage(true);
        clientMessage.setType("jms/text-message");
        // 写入UTF-8编码的字符串长度和内容
        byte[] contentBytes = message.getBytes(StandardCharsets.UTF_8);
        clientMessage.getBodyBuffer().writeInt(contentBytes.length);
        clientMessage.getBodyBuffer().writeBytes(contentBytes);
        producer.send(clientMessage);
    } catch (Exception e) {
        log.error("message = {}", e.getMessage());
    }
}
  • 或者直接使用修复后的JmsTemplate发送消息,自动封装为标准JMS TextMessage,web-console可正常预览。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 02:29:49