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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:37:34