ActiveMQ Artemis无failover协议下独立服务器客户端故障转移实现
ActiveMQ Artemis 手动实现跨独立Broker的客户端故障转移
核心诉求
从ActiveMQ Classic迁移至ActiveMQ Artemis,需手动处理两台独立嵌入式Broker(分属server1、server2)间的客户端故障转移。Artemis不支持Classic的failover:协议,要通过ExceptionListener实现应用层故障转移,达到Classic那样的客户端静默重连效果,且无法配置集群或主备Broker。
背景信息
ActiveMQ Classic客户端配置
import org.apache.activemq.ActiveMQConnectionFactory; @Bean public ConnectionFactory connectionFactory(String user, String credential) { ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory("failover:(tcp://server1:port1,tcp://server2:port2)?jms.useAsyncSend=true&initialReconnectDelay=100&timeout=2000&randomize=false"); factory.setTrustAllPackages(true); factory.setUserName(user); factory.setPassword(credential); return factory; }
ActiveMQ Artemis客户端配置
import org.apache.activemq.artemis.jms.client.ActiveMQJMSConnectionFactory; public ConnectionFactory connectionFactory(String user, String credential) { ActiveMQJMSConnectionFactory factory = new ActiveMQJMSConnectionFactory("(tcp://server1:port1,tcp://server2:port2)?jms.useAsyncSend=true&initialReconnectDelay=100&timeout=2000&sslEnabled=false&failoverAttempts=-1&useTopologyForLoadBalancing=false&connectionLoadBalancingPolicyClassName=org.apache.activemq.artemis.api.core.client.loadbalance.FirstElementConnectionLoadBalancingPolicy", user, credential); return factory; }
现存问题
问题1:初始连接故障转移失败
当server1宕机时,客户端不会自动尝试连接server2,抛出异常:
org.springframework.jms.UncategorizedJmsException: Uncategorized exception occurred during JMS processing; nested exception is ActiveMQNotConnectedException[errorType=NOT_CONNECTED message=AMQ219007: Cannot connect to server(s). Tried with all available servers.]
问题2:已连接后的断开重连失效
客户端连接server1后,若server1关机,客户端不会尝试切换到server2,且ExceptionListener仅在server1重启时才触发。相关日志:
2024-07-29 10:12:45,318 [Thread-0 (ActiveMQ-client-global-threads)] INFO o.a.a.a.c.client [] - AMQ214036: Connection closure to server1/server1:port1 has been detected: AMQ219015: The connection was disconnected because of server shutdown [code=DISCONNECTED] 2024-07-29 10:13:09,546 [Thread-6] INFO d.d.m.j.ExampleListener [] - received exception in topic listener: {} jakarta.jms.JMSException: ActiveMQDisconnectedException[errorType=DISCONNECTED message=AMQ219015: The connection was disconnected because of server shutdown] 2024-07-29 11:32:29,339 [Thread-6] INFO d.d.m.j.ExampleListener [] - received exception in topic listener: {} jakarta.jms.JMSException: ActiveMQDisconnectedException[errorType=DISCONNECTED message=AMQ219015: The connection was disconnected because of server shutdown] at org.apache.activemq.artemis.jms.client.ActiveMQConnection$JMSFailureListener.connectionFailed(ActiveMQConnection.java:724) ~[artemis-jakarta-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.jms.client.ActiveMQConnection$JMSFailureListener.connectionFailed(ActiveMQConnection.java:745) ~[artemis-jakarta-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl.callSessionFailureListeners(ClientSessionFactoryImpl.java:878) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl.callSessionFailureListeners(ClientSessionFactoryImpl.java:866) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl.failoverOrReconnect(ClientSessionFactoryImpl.java:812) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl.handleConnectionFailure(ClientSessionFactoryImpl.java:576) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl$DelegatingFailureListener.connectionFailed(ClientSessionFactoryImpl.java:1417) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.spi.core.protocol.AbstractRemotingConnection.callFailureListeners(AbstractRemotingConnection.java:98) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.protocol.core.impl.RemotingConnectionImpl.fail(RemotingConnectionImpl.java:209) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.core.client.impl.ClientSessionFactoryImpl$CloseRunnable.run(ClientSessionFactoryImpl.java:1182) ~[artemis-core-client-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.utils.actors.OrderedExecutor.doTask(OrderedExecutor.java:57) ~[artemis-commons-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.utils.actors.OrderedExecutor.doTask(OrderedExecutor.java:32) ~[artemis-commons-2.33.0.jar:2.33.0] at org.apache.activemq.artemis.utils.actors.ProcessorBase.executePendingTasks(ProcessorBase.java:68) ~[artemis-commons-2.33.0.jar:2.33.0] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[?:?] at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[?:?] at org.apache.activemq.artemis.utils.ActiveMQThreadFactory$1.run(ActiveMQThreadFactory.java:118) ~[artemis-commons-2.33.0.jar:2.33.0] Caused by: org.apache.activemq.artemis.api.core.ActiveMQDisconnectedException: AMQ219015: The connection was disconnected because of server shutdown ... 7 more 2024-07-29 11:32:29,351 [Thread-6] WARN d.d.m.j.ExampleListener [] - Exception is not an instance of ActiveMQDisconnectedException
当前ExceptionListener实现:
public class ExampleListener implements ExceptionListener { private static final Logger LOG = LoggerFactory.getLogger(ExampleListener.class); @Override public void onException(final JMSException exception) { LOG.info("Received exception in topic listener: {}", exception); if (exception.getLinkedException() instanceof ActiveMQDisconnectedException) { // This if-Statement is always false. ActiveMQDisconnectedException e = (ActiveMQDisconnectedException) exception.getLinkedException(); LOG.info("Disconnected exception type: {}", e.getType()); LOG.info("Disconnected message type: {}", e.getMessage()); // Handle reconnection logic here } else { LOG.warn("Exception is not an instance of ActiveMQDisconnectedException"); } } }
技术疑问
- 若
ExceptionListener仅在服务器重启后才被调用,如何实现主动重连? - 如何编程检测故障Broker并切换连接到另一台,模拟Classic的
failover:协议行为?
解决方案
1. 修复ExceptionListener的异常判断逻辑
Artemis抛出的JMSException中,真实的连接故障异常存储在getCause()中,而非getLinkedException()。修正后的监听器:
public class ExampleListener implements ExceptionListener { private static final Logger LOG = LoggerFactory.getLogger(ExampleListener.class); private final JmsConnectionManager connectionManager; public ExampleListener(JmsConnectionManager connectionManager) { this.connectionManager = connectionManager; } @Override public void onException(final JMSException exception) { LOG.info("Received exception in topic listener: {}", exception); if (exception.getCause() instanceof ActiveMQDisconnectedException || exception.getCause() instanceof ActiveMQNotConnectedException) { LOG.info("Detected connection failure, triggering reconnection..."); connectionManager.reconnect(); } else { LOG.warn("Non-connection related exception occurred", exception); } } }
2. 实现自定义连接管理器(核心重连逻辑)
创建JmsConnectionManager类,负责管理连接、会话、消费者的生命周期,故障时自动切换Broker地址并重建资源:
@Component public class JmsConnectionManager { private static final Logger LOG = LoggerFactory.getLogger(JmsConnectionManager.class); private final List<String> brokerUrls = Arrays.asList("tcp://server1:port1", "tcp://server2:port2"); private int currentBrokerIndex = 0; private Connection connection; private Session session; private MessageConsumer consumer; private final String user; private final String credential; private final MessageListener yourMessageListener; // 替换为业务的消息监听器 @Autowired public JmsConnectionManager(@Value("${jms.user}") String user, @Value("${jms.credential}") String credential, MessageListener yourMessageListener) { this.user = user; this.credential = credential; this.yourMessageListener = yourMessageListener; initializeConnection(); } // 初始化/重建连接 public synchronized void initializeConnection() { try { // 关闭旧资源 closeResources(); // 切换到下一个Broker currentBrokerIndex = (currentBrokerIndex + 1) % brokerUrls.size(); String targetUrl = brokerUrls.get(currentBrokerIndex); LOG.info("Attempting connection to broker: {}", targetUrl); // 创建新连接 ActiveMQJMSConnectionFactory factory = new ActiveMQJMSConnectionFactory(targetUrl, user, credential); connection = factory.createConnection(); connection.setExceptionListener(new ExampleListener(this)); connection.start(); // 重建会话和消费者 session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); Destination destination = session.createTopic("your-topic-name"); // 替换为实际目标 consumer = session.createConsumer(destination); consumer.setMessageListener(yourMessageListener); LOG.info("Successfully connected to broker: {}", targetUrl); } catch (JMSException e) { LOG.error("Failed to connect to broker, retrying in 5 seconds...", e); // 延迟重试,避免频繁请求 Executors.newSingleThreadScheduledExecutor().schedule(this::initializeConnection, 5, TimeUnit.SECONDS); } } // 关闭旧连接资源 private void closeResources() { try { if (consumer != null) consumer.close(); if (session != null) session.close(); if (connection != null) connection.close(); } catch (JMSException e) { LOG.warn("Error closing old JMS resources", e); } } // 对外暴露的重连方法 public void reconnect() { Executors.newSingleThreadExecutor().execute(this::initializeConnection); } }
3. 初始连接故障处理
在Spring Boot启动时,添加初始连接校验,失败则触发重连:
@Configuration public class JmsConfig { private static final Logger LOG = LoggerFactory.getLogger(JmsConfig.class); @Autowired private JmsConnectionManager connectionManager; @PostConstruct public void validateInitialConnection() { try { connectionManager.initializeConnection(); } catch (Exception e) { LOG.error("Initial connection failed, triggering reconnection...", e); connectionManager.reconnect(); } } }
4. 关键配置说明
- 移除Artemis连接URL中的多Broker配置,改用自定义管理器手动切换地址,避免客户端内置负载均衡逻辑干扰。
- 重连逻辑添加延迟重试,防止短时间内频繁重试耗尽资源。
- 使用
synchronized保证连接重建的线程安全,避免并发问题。
内容的提问来源于stack exchange,提问作者LenC
相关产品推荐
相关产品推荐

