如何监听Apache ActiveMQ Broker的连接创建与关闭事件
监控Apache ActiveMQ Broker连接状态的实现方案
方法1:自定义Broker插件拦截连接事件
ActiveMQ支持通过Broker插件拦截Broker层面的连接创建/销毁事件,是监控全局连接状态的最优方案:
- 实现自定义插件,继承
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); } }
- 在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
相关产品推荐
相关产品推荐

