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

Wildfly集群JMS消息驱动Bean跨节点收消息异常及实现咨询

Wildfly集群中JMS/ActiveMQ发布订阅实现方案

问题场景

已搭建Wildfly 23双节点集群(node-1、node-2),使用standalone-full-ha配置文件,集群节点间通信正常。向JMS主题发送消息时,仅发送消息的节点上的MDB能消费,另一节点无法接收,尝试过远程连接工厂、持久化订阅等配置但未解决。

解决方案

1. ActiveMQ Artemis集群配置修改(standalone-full-ha.xml)

需确保消息子系统配置集群发现、集群连接,以及集群感知的连接工厂和主题:

<subsystem xmlns="urn:jboss:domain:messaging-activemq:13.0">
    <server name="default">
        <!-- 广播组:节点间发布自身信息 -->
        <broadcast-group name="bg-group1" jgroups-channel="activemq-cluster">
            <broadcast-period>2000</broadcast-period>
            <connector-ref>http-connector</connector-ref>
        </broadcast-group>
        <!-- 发现组:接收其他节点的广播信息 -->
        <discovery-group name="dg-group1" jgroups-channel="activemq-cluster"/>
        <!-- 集群连接:实现节点间消息路由 -->
        <cluster-connection name="my-cluster" address="jms" connector-name="http-connector" discovery-group="dg-group1"/>
        
        <!-- JMS主题配置:添加对外暴露的条目,确保集群可见 -->
        <jms-topic name="myTopic" entries="java:/jms/topic/myTopic java:jboss/exported/jms/topic/myTopic"/>
        
        <!-- 集群感知连接工厂:用于跨节点发送/接收消息 -->
        <connection-factory name="ClusterConnectionFactory" 
            entries="java:/jms/ClusterConnectionFactory java:jboss/exported/jms/ClusterConnectionFactory" 
            connectors="http-connector" discovery-group="dg-group1"/>
        <!-- 保留原有默认连接工厂,可按需调整 -->
        <connection-factory name="InVmConnectionFactory" entries="java:/ConnectionFactory" connectors="in-vm"/>
        ...
    </server>
</subsystem>

注意:两个节点的配置需保持一致,确保jgroups-channel="activemq-cluster"的通信正常(Wildfly standalone-full-ha.xml默认已配置该JGroups通道)。

2. MDB代码调整

使用持久化订阅,并通过节点名称保证每个节点的订阅唯一:

import javax.ejb.ActivationConfigProperty;
import javax.ejb.MessageDriven;
import javax.jms.Message;
import javax.jms.MessageListener;

@MessageDriven(activationConfig = {
        @ActivationConfigProperty(propertyName = "destinationType", propertyValue = "javax.jms.Topic"),
        @ActivationConfigProperty(propertyName = "destinationLookup", propertyValue = "java:/jms/topic/myTopic"),
        @ActivationConfigProperty(propertyName = "maxSession", propertyValue = "1"),
        // 启用持久化订阅,保证集群节点都能收到消息
        @ActivationConfigProperty(propertyName = "subscriptionDurability", propertyValue = "Durable"),
        // 用节点名称作为clientId后缀,避免集群内clientId冲突
        @ActivationConfigProperty(propertyName = "clientId", propertyValue = "ClusteredEventListener-${jboss.node.name}"),
        // 唯一订阅名称,每个节点对应独立订阅
        @ActivationConfigProperty(propertyName = "subscriptionName", propertyValue = "ClusteredSubscription-${jboss.node.name}")
})
public class ClusteredEventListener implements MessageListener {

    @Override
    public void onMessage(final Message message) {
        try {
            // 示例:打印消费节点及消息ID
            String nodeName = System.getProperty("jboss.node.name");
            System.out.printf("Node [%s] consumed message, ID: %s%n", nodeName, message.getJMSMessageID());
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

${jboss.node.name}会被Wildfly自动替换为当前节点名称(如node-1),确保每个节点的订阅标识唯一。

3. 消息发送代码优化

使用集群连接工厂,并启动连接以确保消息路由到集群:

import java.io.Serializable;
import java.util.logging.Logger;

import javax.annotation.Resource;
import javax.ejb.Startup;
import javax.enterprise.context.ApplicationScoped;
import javax.jms.JMSException;
import javax.jms.MessageProducer;
import javax.jms.ObjectMessage;
import javax.jms.Session;
import javax.jms.Topic;
import javax.jms.TopicConnection;
import javax.jms.TopicConnectionFactory;
import javax.jms.TopicSession;

@Startup
@ApplicationScoped
public class ClusteredEventSender {
    private static final Logger log = Logger.getLogger(ClusteredEventSender.class.getName());

    @Resource(lookup = "java:/jms/topic/myTopic")
    private Topic topic;

    // 注入集群连接工厂
    @Resource(lookup = "java:/jms/ClusterConnectionFactory")
    private TopicConnectionFactory connectionFactory;

    public void broadcast(final Serializable event) {
        try (TopicConnection connection = connectionFactory.createTopicConnection()) {
            connection.start(); // 必须启动连接,否则消息无法路由到集群节点
            try (TopicSession session = connection.createTopicSession(false, Session.AUTO_ACKNOWLEDGE)) {
                try (MessageProducer messageProducer = session.createPublisher(topic)) {
                    // 按需设置消息持久化模式,持久化消息会在节点重启后保留
                    messageProducer.setDeliveryMode(javax.jms.DeliveryMode.NON_PERSISTENT);
                    ObjectMessage message = session.createObjectMessage(event);
                    messageProducer.send(message);
                    log.info("Event broadcasted to cluster topic: " + event.toString());
                }
            }
        } catch (JMSException e) {
            log.severe("Failed to broadcast event: " + event.toString() + ", error: " + e.getMessage());
        }
    }
}

验证步骤

  1. 将修改后的配置同步到两个集群节点
  2. 部署调整后的MDB和消息发送者到两个节点
  3. 调用node-1上的ClusteredEventSender.broadcast()方法发送消息
  4. 查看两个节点的日志,确认均输出消息消费记录

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 18:40:29