RabbitMQ客户端Connection复用
方案概述
在RabbitMQ中,物理TCP连接(Connection)是一种昂贵且受限的资源,每条连接在Broker端至少占用约100KB内存。大量连接会触发实例流控或达到实例连接数上限,导致新连接被拒绝。合理复用连接是保障RabbitMQ实例稳定运行的关键。
本章节针对如表1 推荐方案中的场景,分别提供不同的Connection复用方案。
| 场景 | 推荐方案 |
|---|---|
| 通用客户端(Java/Go/Python/.NET/Node.js)场景 | 遵循通用客户端场景方案,复用Connection+Channel。 |
| PHP-FPM客户端短连接场景 | 通过部署AMQProxy到PHP短连接场景,复用Connection,无需改造业务代码。 |
通用客户端场景
| 原则 | 说明 |
|---|---|
| Connection必须复用 | 应用启动时创建,整个生命周期重复使用,严禁每次发送消息都新建。 |
| Channel不得跨线程共享 | 每个线程使用独立Channel。 |
| 生产与消费隔离 | 生产与消费使用不同的Connection,避免消费端流控影响生产端。 |
| 开启自动恢复 | 应对Broker重启、网络闪断、主备切换。 |
| 心跳超时时间设置为10秒 | RabbitMQ实例LVS心跳超时时间为90秒,客户端心跳超时时间必须小于90秒。 |
| 配置项 | 推荐值 |
|---|---|
| setAutomaticRecoveryEnabled(true) | true |
| setTopologyRecoveryEnabled(true) | true |
| setNetworkRecoveryInterval(5000) | 5000,单位:ms |
| setRequestedHeartbeat(10) | 10,单位:s |
| setConnectionTimeout(5000) | 5000,单位:ms |
Connection管理(双重检查锁)代码示例:
public class RabbitMQConnectionManager {
private static volatile Connection connection;
private static final Object lock = new Object();
public static Connection getConnection() throws Exception {
Connection local = connection;
if (local == null || !local.isOpen()) {
synchronized (lock) {
local = connection;
if (local == null || !local.isOpen()) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost(System.getenv("MQ_HOST"));
factory.setPort(5672);
factory.setUsername(System.getenv("MQ_USER"));
factory.setPassword(System.getenv("MQ_PASSWORD"));
factory.setVirtualHost("/");
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
factory.setNetworkRecoveryInterval(5000);
factory.setRequestedHeartbeat(10);
factory.setConnectionTimeout(5000);
local = factory.newConnection();
connection = local;
}
}
}
return local;
}
} Channel使用代码示例:
try (Channel channel = connection.createChannel()) {
channel.basicPublish(exchange, routingKey, null, message.getBytes());
} 部署AMQProxy到PHP短连接场景
PHP-FPM每个请求独立且短暂,无法维持长连接。每次请求新建Connection会导致AMQP握手开销大、连接数飙升、Channel ID耗尽。
本方案将AMQProxy部署在客户端与RabbitMQ实例之间,将客户端短连接转换为AMQProxy与RabbitMQ实例之间的少量长连接。开启连接时,按用户名+密码+Vhost复用上游Connection,未命中则新建Connection,关闭连接时直接应答OK,不真正关闭上游连接,留待复用。
部署方式如下:
- 下载AMQProxy,并解压到本地目录。本文以v3.1.2为例,也可以下载GitHub Releases其他版本。
wget https://github.com/cloudamqp/amqproxy/releases/download/v3.1.2/amqproxy-3.1.2_static-amd64.tar.gz tar -xzvf amqproxy-3.1.2_static-amd64.tar.gz
- 进入解压后的AMQProxy文件夹路径,启动AMQProxy。
cd amqproxy/ ./amqproxy -l 127.0.0.1 -p 5673 amqp://rabbitmq.example.com:5672
其中,rabbitmq.example.com:5672为RabbitMQ实例的连接地址和端口,请替换为RabbitMQ实际的连接地址和端口。
启动成功输出:
amq_proxy.server Proxy upstream: rabbitmq.example.com:5672 amq_proxy.http_server Bound to 127.0.0.1:15673 amq_proxy.http_server HTTP server listening on 127.0.0.1:15673 amq_proxy.server Proxy listening on 127.0.0.1:5673
生产环境建议使用systemd守护进程,并部署多实例+负载均衡避免单点故障。
- 在RabbitMQ客户端将连接地址改为AMQProxy的IP与端口,用户名、密码、Vhost保持不变。
// 直连RabbitMQ: // $conn = new AMQPStreamConnection('rabbitmq.example.com', 5672, $user, $pass, '/'); // 经AMQProxy接入: $conn = new AMQPStreamConnection('127.0.0.1', 5673, $user, $pass, '/'); - 通过以下几点进行验证。
- AMQProxy输出Proxy listening on即监听成功。
- 在RabbitMQ控制台查看“连接数”监控,确认连接数显著下降。查看监控的操作请参见查看RabbitMQ监控数据。