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"的通信正常(Wildflystandalone-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()); } } }
验证步骤
- 将修改后的配置同步到两个集群节点
- 部署调整后的MDB和消息发送者到两个节点
- 调用node-1上的
ClusteredEventSender.broadcast()方法发送消息 - 查看两个节点的日志,确认均输出消息消费记录
内容的提问来源于stack exchange,提问作者Lars Gendner
相关产品推荐
相关产品推荐

