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

如何监听Apache ActiveMQ Broker的连接创建与关闭事件

监控Apache ActiveMQ Broker连接状态的实现方案

方法1:自定义Broker插件拦截连接事件

ActiveMQ支持通过Broker插件拦截Broker层面的连接创建/销毁事件,是监控全局连接状态的最优方案:

  1. 实现自定义插件,继承BrokerFilter并重写连接相关方法:
import org.apache.activemq.broker.Broker;
import org.apache.activemq.broker.BrokerFilter;
import org.apache.activemq.broker.BrokerPlugin;
import org.apache.activemq.command.ConnectionInfo;

public class ConnectionMonitorPlugin extends BrokerFilter implements BrokerPlugin {

    public ConnectionMonitorPlugin(Broker next) {
        super(next);
    }

    @Override
    public void addConnection(ConnectionInfo info) throws Exception {
        // 新连接创建时执行自定义逻辑,比如发送通知、记录日志
        System.out.println("新连接建立: 客户端ID=" + info.getClientId() + ",远程地址=" + info.getRemoteAddress());
        super.addConnection(info);
    }

    @Override
    public void removeConnection(ConnectionInfo info, Throwable error) throws Exception {
        // 连接关闭时执行自定义逻辑
        System.out.println("连接关闭: 客户端ID=" + info.getClientId() + ",远程地址=" + info.getRemoteAddress());
        super.removeConnection(info, error);
    }

    @Override
    public Broker installPlugin(Broker broker) throws Exception {
        return new ConnectionMonitorPlugin(broker);
    }
}
  1. 在ActiveMQ配置文件activemq.xml中注册插件:
<broker xmlns="http://activemq.apache.org/schema/core" brokerName="localhost" dataDirectory="${activemq.data}">
    <!-- 其他原有配置 -->
    <plugins>
        <bean xmlns="http://www.springframework.org/schema/beans" class="com.your.package.ConnectionMonitorPlugin"/>
    </plugins>
    <!-- 其他原有配置 -->
</broker>

方法2:通过JMX监控全局连接

ActiveMQ默认暴露JMX MBean,可通过JMX连接监听连接事件:

import javax.management.MBeanServerConnection;
import javax.management.Notification;
import javax.management.NotificationListener;
import javax.management.ObjectName;
import javax.management.remote.JMXConnector;
import javax.management.remote.JMXConnectorFactory;
import javax.management.remote.JMXServiceURL;

public class JmxConnectionMonitor {
    public static void main(String[] args) throws Exception {
        // 连接ActiveMQ默认JMX地址
        JMXServiceURL jmxUrl = new JMXServiceURL("service:jmx:rmi:///jndi/rmi://localhost:1099/jmxrmi");
        JMXConnector connector = JMXConnectorFactory.connect(jmxUrl);
        MBeanServerConnection mBeanConn = connector.getMBeanServerConnection();

        // 监听Broker的连接通知
        ObjectName brokerMBean = new ObjectName("org.apache.activemq:type=Broker,brokerName=localhost");
        mBeanConn.addNotificationListener(brokerMBean, (Notification notification, Object handback) -> {
            switch (notification.getType()) {
                case "ConnectionAdded":
                    System.out.println("新连接加入: " + notification.getMessage());
                    break;
                case "ConnectionRemoved":
                    System.out.println("连接移除: " + notification.getMessage());
                    break;
            }
        }, null, null);
    }
}

注:确保ActiveMQ已开启JMX(默认开启),若需调整JMX端口等参数,可修改activemq.xml或activemq.env配置。

方法3:客户端侧监控自身连接

如果只需监控当前客户端与Broker的连接状态,可使用ConnectionListener:

import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.ConnectionListener;

public class ClientConnectionWatcher {
    public static void main(String[] args) throws Exception {
        ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616");
        ActiveMQConnection connection = (ActiveMQConnection) factory.createConnection();

        connection.addConnectionListener(new ConnectionListener() {
            @Override
            public void onConnectionStarted() {
                System.out.println("客户端连接启动");
            }

            @Override
            public void onConnectionStopped() {
                System.out.println("客户端连接停止");
            }

            @Override
            public void onConnectionInterrupted() {
                System.out.println("客户端连接中断");
            }

            @Override
            public void onConnectionResumed() {
                System.out.println("客户端连接恢复");
            }

            @Override
            public void onConnectionException(Throwable exception) {
                System.err.println("客户端连接异常: " + exception.getMessage());
            }
        });

        connection.start();
        // 业务逻辑代码
    }
}

该方法仅能监控当前客户端自身的连接状态,无法感知其他客户端的连接变化。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 00:10:52