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

ActiveMQ Artemis集群大消息重分发失效、消息阻塞问题

ActiveMQ Artemis 集群大消息投递阻塞问题

环境说明

使用ActiveMQ Artemis 2.23.1搭建3主3从复制模式集群,所有节点均配置60GB硬盘、60GB内存,客户端接收约100MB的大消息时出现投递异常。

测试流程

  • 在node01、node03节点部署消费者
  • 发送100条小消息
  • 发送100条100MB级别的大消息
  • 再次发送100条小消息

异常表现

  • 首批100条小消息可正常投递,发送大消息后Broker直接进入阻塞状态:大消息始终无法被消费者接收,最后一批次发送的小消息也无法正常投递
  • 监控可见$.artemis.internal.sf.amq-cluster.<id>格式的集群内部队列存在大消息积压,即使队列上存在消费者也无法消费这些积压消息
  • 问题可稳定复现:基于官方源码中features > clustered > queue-message-redistribution示例修改,集成features > standard > large-messages示例的大消息收发逻辑,使用2个本地嵌入式服务执行mvn verify命令即可100%复现阻塞现象

版本回归测试结果

2022年7月25日补充测试数据:分别在2.19.1(Java 8运行环境)、2.21(Java 11运行环境)、2.22(Java 11运行环境)版本的嵌入式环境运行测试用例,确认该阻塞问题从2.22版本开始引入,后续将在2.21版本的类生产环境中开展进一步验证。

现有集群配置

所有集群配置文件均通过Ansible脚本统一生成,怀疑遗漏大消息处理相关的关键配置项,各配置文件内容如下:

broker.xml

<?xml version='1.0'?>
<configuration xmlns="urn:activemq"
               xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
               xmlns:xi="http://www.w3.org/2001/XInclude"
               xsi:schemaLocation="urn:activemq /schema/artemis-configuration.xsd">

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

        <name>master01.intra</name>

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

        <journal-type>ASYNCIO</journal-type>

        <paging-directory>data/paging</paging-directory>

        <bindings-directory>data/bindings</bindings-directory>

        <journal-directory>data/journal</journal-directory>

        <large-messages-directory>data/large-messages</large-messages-directory>

        <journal-datasync>true</journal-datasync>
        <!--        <journal-sync-non-transactional>false</journal-sync-non-transactional>-->
        <!--        <journal-sync-transactional>false</journal-sync-transactional>-->
        <journal-min-files>2</journal-min-files>

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

        <journal-buffer-timeout>280000</journal-buffer-timeout>


        <!-- how often we are looking for how many bytes are being used on the disk in ms -->
        <disk-scan-period>5000</disk-scan-period>

        <!-- once the disk hits this limit the system will block, or close the connection in certain protocols
             that won't support flow control. -->
        <max-disk-usage>90</max-disk-usage>

        <!-- should the broker detect dead locks and other issues -->
        <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>552000</page-sync-timeout>

        <connectors>
            <!-- Connector used to be announced through cluster connections and notifications -->
            <connector name="artemis">tcp://master01.intra:61616</connector>
        </connectors>

        <acceptors>
            <!-- useEpoll means: it will use Netty epoll if you are on a system (Linux) that supports it -->
            <!-- amqpCredits: The number of credits sent to AMQP producers -->
            <!-- amqpLowCredits: The server will send the # credits specified at amqpCredits at this low mark -->
                        <acceptor name="artemis">tcp://master01.intra:61616?tcpSendBufferSize=1048576;tcpReceiveBufferSize=1048576;protocols=CORE,OPENWIRE;useEpoll=true;amqpCredits=1000;amqpLowCredits=300;virtualTopicConsumerWildcards=Consumer.*.%3E%3B2;supportAdvisory=false;suppressInternalManagementObjects=false</acceptor>
                    </acceptors>

        <cluster-user>admin-cluster</cluster-user>
        <cluster-password>admin-cluster</cluster-password>

        <broadcast-groups>
            <broadcast-group name="bg-group1">
                <group-address>231.7.7.7</group-address>
                <group-port>9876</group-port>
                <broadcast-period>2000</broadcast-period>
                <connector-ref>artemis</connector-ref>
            </broadcast-group>
        </broadcast-groups>

        <discovery-groups>
            <discovery-group name="dg-group1">
                <group-address>231.7.7.7</group-address>
                <group-port>9876</group-port>
                <refresh-timeout>10000</refresh-timeout>
            </discovery-group>
        </discovery-groups>

        <cluster-connections>
            <cluster-connection name="amq-cluster">
                <address></address>
                <connector-ref>artemis</connector-ref>
                <message-load-balancing>ON_DEMAND</message-load-balancing>
                <discovery-group-ref discovery-group-name="dg-group1"/>
            </cluster-connection>
        </cluster-connections>

        <xi:include href="/app/esbbroker/etc/security-settings.xml"/>
        <xi:include href="/app/esbbroker/etc/addresses-settings.xml"/>
        <xi:include href="/app/esbbroker/etc/addresses.xml"/>
        <xi:include href="/app/esbbroker/etc/ha-policy.xml"/>

        <metrics>
            <jvm-memory>true</jvm-memory> <!-- defaults to true -->
            <jvm-gc>true</jvm-gc> <!-- defaults to false -->
            <jvm-threads>true</jvm-threads> <!-- defaults to false -->
            <netty-pool>false</netty-pool> <!-- defaults to false -->
            <plugin class-name="org.apache.activemq.artemis.core.server.metrics.plugins.ArtemisPrometheusMetricsPlugin"/>
        </metrics>
    </core>
</configuration>

addresses.xml

该文件包含大量地址配置,所有配置格式统一如下:

<addresses xmlns="urn:activemq:core">
    <address name="stirint.clo.person.signal">
        <anycast>
            <queue name="stirint.clo.person.signal"/>
        </anycast>
    </address>
    ...
</addresses>

addresses-settings.xml

全局匹配规则match="#"中redistribution-delay配置为0,与官方queue-redistribution示例配置保持一致

<!-- example of one of the clients queue (all generated the same way by ansible -->
<address-settings xmlns="urn:activemq:core">
    <address-setting match="stirint.clo.person.signal">
        <dead-letter-address>DLQ.stirint.clo.person.signal</dead-letter-address>
        <auto-create-dead-letter-resources>true</auto-create-dead-letter-resources>
        <max-delivery-attempts>3</max-delivery-attempts>
        <expiry-address>ExpiryQueue</expiry-address>
        <redelivery-delay>0</redelivery-delay>
        <!-- with -1 only the global-max-size is in use for limiting -->
        <max-size-bytes>-1</max-size-bytes>
        <message-counter-history-day-limit>10</message-counter-history-day-limit>
        <address-full-policy>PAGE</address-full-policy>
        <auto-create-queues>true</auto-create-queues>
        <auto-create-addresses>true</auto-create-addresses>
        <auto-delete-queues>false</auto-delete-queues>
        <auto-delete-addresses>false</auto-delete-addresses>
        <auto-create-jms-queues>false</auto-create-jms-queues>
        <auto-create-jms-topics>false</auto-create-jms-topics>
    </address-setting>
    <!--  ... -->
    <!-- other entries -->
    <!-- if you define auto-create on certain queues, management has to be auto-create -->
    <address-setting match="activemq.management#">
        <dead-letter-address>DLQ</dead-letter-address>
        <expiry-address>ExpiryQueue</expiry-address>
        <redelivery-delay>0</redelivery-delay>
        <!-- with -1 only the global-max-size is in use for limiting -->
        <max-size-bytes>-1</max-size-bytes>
        <message-counter-history-day-limit>10</message-counter-history-day-limit>
        <address-full-policy>PAGE</address-full-policy>
        <auto-create-queues>true</auto-create-queues>
        <auto-create-addresses>true</auto-create-addresses>
        <auto-delete-queues>false</auto-delete-queues>
        <auto-delete-addresses>false</auto-delete-addresses>
        <auto-create-jms-queues>false</auto-create-jms-queues>
        <auto-create-jms-topics>false</auto-create-jms-topics>
    </address-setting>
    <!--default for catch all-->
    <address-setting match="#">
        <dead-letter-address>DLQ</dead-letter-address>
        <expiry-address>ExpiryQueue</expiry-address>
        <redelivery-delay>0</redelivery-delay>
        <redistribution-delay>0</redistribution-delay>
        <!-- with -1 only the global-max-size is in use for limiting -->
        <max-size-bytes>-1</max-size-bytes>
        <message-counter-history-day-limit>10</message-counter-history-day-limit>
        <address-full-policy>PAGE</address-full-policy>
        <auto-create-queues>true</auto-create-queues>
        <auto-create-addresses>true</auto-create-addresses>
        <auto-delete-queues>false</auto-delete-queues>
        <auto-delete-addresses>false</auto-delete-addresses>
        <auto-create-jms-queues>false</auto-create-jms-queues>
        <auto-create-jms-topics>false</auto-create-jms-topics>
    </address-setting>
</address-settings>

security-settings.xml

<security-settings xmlns="urn:activemq:core">
    <!-- one of many, all have similar format because generated by ansible script -->
    <security-setting match="stirint.clo.person.signal">
        <permission type="consume" roles="gStirint,amq"/>
        <permission type="browse" roles="gStirint,amq,readonly"/>
        <permission type="send" roles="gStirint,amq"/>
        <permission type="createNonDurableQueue" roles="gStirint,amq"/>
        <permission type="deleteNonDurableQueue" roles="gStirint,amq"/>
        <permission type="createDurableQueue" roles="gStirint,amq"/>
        <permission type="deleteDurableQueue" roles="gStirint,amq"/>
        <permission type="createAddress" roles="gStirint,amq"/>
        <permission type="deleteAddress" roles="gStirint,amq"/>
    </security-setting>
    ...
    <security-setting match="ActiveMQ.Advisory.TempQueue">
        <permission type="createNonDurableQueue" roles="amq,readonly" />
        <permission type="deleteNonDurableQueue" roles="amq,readonly" />
        <permission type="createDurableQueue" roles="amq,readonly" />
        <permission type="browse" roles="amq,readonly"/>
        <permission type="send" roles="amq,readonly"/>
    </security-setting>

    <security-setting match="ActiveMQ.Advisory.TempTopic">
        <permission type="createNonDurableQueue" roles="amq,readonly"/>
        <permission type="deleteNonDurableQueue" roles="amq,readonly"/>
        <permission type="createDurableQueue" roles="amq,readonly" />
        <permission type="browse" roles="amq,readonly"/>
        <permission type="send" roles="amq,readonly"/>
    </security-setting>

    <security-setting match="#">
        <permission type="createNonDurableQueue" roles="amq"/>
        <permission type="deleteNonDurableQueue" roles="amq"/>
        <permission type="createDurableQueue" roles="amq"/>
        <permission type="deleteDurableQueue" roles="amq"/>
        <permission type="createAddress" roles="amq"/>
        <permission type="deleteAddress" roles="amq"/>
        <permission type="consume" roles="amq"/>
        <permission type="browse" roles="amq,readonly"/>
        <permission type="send" roles="amq"/>
        <!-- we need this otherwise ./artemis data imp wouldn't work -->
        <permission type="manage" roles="amq"/>
    </security-setting>
</security-settings>

ha-policy.xml

集群共3组主从节点,分属gn-1、gn-2、gn-3三个分组,单节点配置示例如下:

<ha-policy xmlns="urn:activemq:core">
    <!-- this file contains configuration of high availability cluster -->
    <replication>
        <master>
            <check-for-live-server>true</check-for-live-server>
            <group-name>gn-1</group-name>
        </master>
    </replication>
</ha-policy>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 16:15:37