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

使用嵌入式ActiveMQ Artemis实现发布订阅时出现轮询分发问题

问题原因

你当前的实现是点对点模式,而非发布订阅模式。所有消费者都绑定到同一个队列q1,而ActiveMQ Artemis中队列的默认消息分发策略是轮询(round-robin),因此消息会依次发给不同消费者,这是队列的标准点对点行为。发布订阅模式需要让每个订阅者拥有独立的消息副本,需利用地址的多播特性,让每个订阅者对应专属队列(临时或持久化)。


解决方案

方案1:自动创建临时订阅队列(推荐,适合非持久化场景)

直接让消费者订阅地址address1,而非指定固定队列。Broker会自动为每个消费者创建临时队列并绑定到地址的多播路由上,每个消息会被发送到所有临时队列,实现所有消费者都收到消息。

修改消费者代码

public void createConsumer(int id) throws Exception {
    ServerLocator locator = ActiveMQClient.createServerLocator("vm://0");
    ClientSessionFactory factory = locator.createSessionFactory();
    ClientSession session = factory.createSession();
    session.start();
    // 直接订阅地址,而非固定队列
    ClientConsumer consumer = session.createConsumer("address1");
    consumer.setMessageHandler(new MessageHandler() {
        @Override
        public void onMessage(ClientMessage clientMessage) {
            log_.info(">>> consumer {} received message content: {}", id, clientMessage.getBodyBuffer().readString());
        }
    });
}

简化broker.xml配置

删除地址下的固定队列,仅保留多播配置:

<addresses>
    <address name="address1">
        <multicast/> <!-- 无需固定队列,临时队列自动创建 -->
    </address>
</addresses>

方案2:手动创建多个固定队列(适合持久化或自定义队列属性场景)

如果需要为每个消费者指定固定队列,可在地址的多播节点下创建多个队列,让每个消费者绑定到不同队列。

修改broker.xml配置

<addresses>
    <address name="address1">
        <multicast>
            <queue name="q1" max-consumers="99">
                <durable>false</durable>
            </queue>
            <queue name="q2" max-consumers="99">
                <durable>false</durable>
            </queue>
            <queue name="q3" max-consumers="99">
                <durable>false</durable>
            </queue>
        </multicast>
    </address>
</addresses>

修改消费者代码

让每个消费者绑定到对应固定队列:

public void createConsumer(int id) throws Exception {
    ServerLocator locator = ActiveMQClient.createServerLocator("vm://0");
    ClientSessionFactory factory = locator.createSessionFactory();
    ClientSession session = factory.createSession();
    session.start();
    // 每个消费者绑定到不同的固定队列
    ClientConsumer consumer = session.createConsumer("q" + (id + 1));
    consumer.setMessageHandler(new MessageHandler() {
        @Override
        public void onMessage(ClientMessage clientMessage) {
            log_.info(">>> consumer {} received message content: {}", id, clientMessage.getBodyBuffer().readString());
        }
    });
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 15:34:52