Connection Reuse on the RabbitMQ Client
Overview
In RabbitMQ, physical TCP connections are expensive and limited. Each connection occupies at least 100 KB memory on the broker. A large number of connections will trigger instance flow control. If the number of instance connections reaches the upper limit, new connections will be rejected. Proper connection reuse is the key to ensuring the stable running of RabbitMQ instances.
This section provides different connection reuse solutions for the scenarios listed in Table 1.
| Scenario | Recommended Solution |
|---|---|
| Common clients (Java/Go/Python/.NET/Node.js) | Use the solution in Common Client Scenario to reuse connections and channels. |
| PHP-FPM clients (short connections) | Use the solution in Deploying AMQProxy in the PHP Short Connection Scenario to reuse connections, without the need to reconstruct service code. |
Common Client Scenario
| Principle | Description |
|---|---|
| Reusing connections | A connection is created when the application is started and is reused throughout the lifecycle. It is prohibited to create a connection each time a message is sent. |
| Do not share channels across threads | Each thread uses an independent channel. |
| Isolation between production and consumption | Production and consumption use different connections to prevent flow control on the consumer from affecting the producer. |
| Enabling automatic recovery | This is to cope with broker restart, intermittent network disconnection, and active/standby switchover. |
| Setting the heartbeat timeout interval to 10 seconds | The LVS heartbeat timeout interval of a RabbitMQ instance is 90 seconds. The heartbeat timeout interval of the client must be shorter than 90 seconds. |
| Configuration Item | Recommended Value |
|---|---|
| setAutomaticRecoveryEnabled(true) | true |
| setTopologyRecoveryEnabled(true) | true |
| setNetworkRecoveryInterval(5000) | 5000, in ms |
| setRequestedHeartbeat(10) | 10, in s |
| setConnectionTimeout(5000) | 5000, in ms |
Example code for managing connections (double-checked locking):
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;
}
} Example code of using a channel:
try (Channel channel = connection.createChannel()) {
channel.basicPublish(exchange, routingKey, null, message.getBytes());
} Deploying AMQProxy in the PHP Short Connection Scenario
Each PHP-FPM request is independent and short-lived, and persistent connections cannot be maintained. Creating a connection for each request will cause high AMQP handshake overhead, a sharp increase in the number of connections, and exhaustion of channel IDs.
In this solution, deploy AMQProxy between the client and the RabbitMQ instance to convert the client's short connections into a small number of persistent connections between AMQProxy and the RabbitMQ instance. When a connection is enabled, the upstream connection is reused based on the username, password, and virtual host. If the connection is not hit, a new connection will be created. When the connection is disabled, OK is returned. In this case, the upstream connection is not actually disabled; instead, it is reserved for reuse.
Deployment procedure:
- Download AMQProxy and decompress it to a local directory. This document uses v3.1.2 as an example. You can also download other 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
- Go to the directory where the AMQProxy file is decompressed and start AMQProxy.
cd amqproxy/ ./amqproxy -l 127.0.0.1 -p 5673 amqp://rabbitmq.example.com:5672
In the preceding, rabbitmq.example.com:5672 indicates the connection address and port number of the RabbitMQ instance. Replace it with the actual connection address and port number of the RabbitMQ instance.
Output after successful startup:
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
In the production environment, you are advised to use the systemd daemon and deploy multiple instances and load balancing to avoid single points of failure.
- On the RabbitMQ client, change the connection address to the IP address and port number of AMQProxy, and keep the username, password, and virtual host unchanged.
// Directly connect to RabbitMQ: // $conn = new AMQPStreamConnection('rabbitmq.example.com', 5672, $user, $pass, '/'); // Connect to RabbitMQ through AMQProxy: $conn = new AMQPStreamConnection('127.0.0.1', 5673, $user, $pass, '/'); - Implement verification.
- If AMQProxy outputs "Proxy listening on", the listening is successful.
- On the RabbitMQ console, check whether the value of Connections decreases significantly. For details about how to view the monitoring data, see Viewing RabbitMQ Metrics.
What is your overall rating for this page?
Thank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot