amqplib and amqp-client.js, focusing especially on connection handling and recovery.
Now, let’s put those differences into practice by migrating our application from amqplib to the high-level amqp-client.js API.
Because the LavinMQ-specific operations are already isolated in our broker layer, we don’t need to change the producer, worker, or transaction-processing logic. Our main migration is amqplib/broker.js. We will keep the same broker interface and change only the client-specific parts needed to communicate with LavinMQ.
Step 1: Install amqp-client.js
Install the new client:
npm install @cloudamqp/amqp-client
Copy amqplib/broker.js to amqp-client/broker.js and update amqp-client/broker.js. Replace:
const amqp = require("amqplib");
with:
const { AMQPSession } = require("@cloudamqp/amqp-client");
Step 2: Replace the connection and channel
Replace:
connection = await amqp.connect(config.amqpUrl);
channel = await connection.createChannel();
With:
session = await AMQPSession.connect(config.amqpUrl, {
reconnectInterval: 1000,
maxReconnectInterval: 30000,
backoffMultiplier: 2,
maxRetries: 0,
onconnect: () => {
log("Connected to LavinMQ");
},
ondisconnect: (error) => {
log("Disconnected from LavinMQ", {
error: error?.message,
});
},
onfailed: (error) => {
log("Failed to reconnect to LavinMQ", {
error: error?.message,
});
},
});
We no longer need to explicitly create a channel in our broker. We also need references to the queues and exchange that we publish to:
transactionsQueue = await session.queue(
config.queues.transactions,
{
durable: true,
}
);
retryQueue = await session.queue(
config.queues.retry,
{
durable: true,
}
);
deadExchange = await session.directExchange(
config.exchanges.dlx,
{
durable: true,
}
);
These replace the places where our amqplib implementation publishes through the channel.
Step 3: Replace sendToQueue() for transaction publishing
Replace:
channel.sendToQueue(
config.queues.transactions,
Buffer.from(JSON.stringify(transaction)),
{
persistent: true,
contentType: "application/json",
headers: {
scenario,
},
}
);
With:
await transactionsQueue.publish(
JSON.stringify(transaction),
{
deliveryMode: 2,
contentType: "application/json",
headers: {
scenario,
},
}
);
We don’t need to change producer.js. It still calls the same broker function: await broker.publishTransaction(transaction, scenario); Only the implementation behind that function has changed.
Step 4: Replace the consumer
With amqplib, we configure prefetch and consume through the channel:
await channel.prefetch(1);
await channel.consume(
config.queues.transactions,
(message) => {
// handle message
},
{
noAck: false,
}
);
With amqp-client.js, replace that with:
await transactionsQueue.subscribe(
{
prefetch: 1,
manualAck: true,
},
async (message) => {
// handle message
}
);
There is also a small difference in how we read the message body. With amqplib:
const transaction = JSON.parse(
message.content.toString()
);
Replace it with:
const transaction = JSON.parse(
message.bodyString()
);
The broker still passes the same normalized data to our worker, so there is no need to change worker.js.
Step 5: Replace message acknowledgement
With amqplib, acknowledgement goes through the channel:
channel.ack(message);
Replace it with:
await message.ack();
The worker still calls await broker.acknowledge(originalMessage); so the change stays inside the broker.
Step 6: Replace retry publishing
Our retry strategy doesn’t change. We still send retryable transactions to the retry queue with an x-retry-count header. With amqplib:
channel.sendToQueue(
config.queues.retry,
Buffer.from(JSON.stringify(transaction)),
{
persistent: true,
contentType: "application/json",
headers: {
"x-retry-count": retryCount,
},
}
);
Replace it with:
await retryQueue.publish(
JSON.stringify(transaction),
{
deliveryMode: 2,
contentType: "application/json",
headers: {
"x-retry-count": retryCount,
},
}
);
Nothing about the retry strategy itself has changed. We are only changing how the broker publishes the retry message.
Step 7: Replace dead-letter publishing
With amqplib, we publish a permanently failed transaction through the dead-letter exchange using:
channel.publish(
config.exchanges.dlx,
config.queues.dead,
Buffer.from(
JSON.stringify({
transaction,
failedAt: new Date().toISOString(),
reason: error.message,
retries,
})
),
{
persistent: true,
contentType: "application/json",
}
);
Replace it with:
await deadExchange.publish(
JSON.stringify({
transaction,
failedAt: new Date().toISOString(),
reason: error.message,
retries,
}),
{
routingKey: config.queues.dead,
deliveryMode: 2,
contentType: "application/json",
}
);
The dead-letter exchange, routing key, queue, and failure-handling strategy remain the same. Only the publishing API changes.
Step 8: Replace the custom reconnection handling
This is one of the more important changes in our migration. A connection to LavinMQ can disappear for many reasons — a short network interruption, a broker restart, or another interruption. The important question for a long-running worker is:
What happens after the connection is lost?
In our amqplib implementation, we handle that ourselves. We listen for the connection closing:
connection.on("close", reconnect);
Our reconnect() function then waits five seconds and calls connect() again:
function reconnect() {
if (shuttingDown || !messageHandler || reconnectTimer) {
return;
}
log("Connection closed. Reconnecting in 5 seconds...");
reconnectTimer = setTimeout(async () => {
reconnectTimer = null;
try {
await connect();
await subscribe(messageHandler);
log("Reconnected to LavinMQ");
} catch (error) {
log("Reconnect failed", {
error: error.message,
});
reconnect();
}
}, 5000);
}
Notice that reconnecting the TCP/AMQP connection alone isn’t enough for our worker. After creating the new connection, our code also calls:
await subscribe(messageHandler);
This is important because the worker needs to continue consuming messages after the connection comes back. In our amqplib implementation, our application is responsible for the recovery flow. We also keep messageHandler so that we know which handler to register again after reconnecting.
With our high-level amqp-client.js implementation, we don’t copy this custom reconnect() loop. Instead, we configure connection recovery when creating the session:
session = await AMQPSession.connect(config.amqpUrl, {
reconnectInterval: 1000,
maxReconnectInterval: 30000,
backoffMultiplier: 2,
maxRetries: 0,
onconnect: () => {
log("Connected to LavinMQ");
},
ondisconnect: (error) => {
log("Disconnected from LavinMQ", {
error: error?.message,
});
},
onfailed: (error) => {
log("Failed to reconnect to LavinMQ", {
error: error?.message,
});
},
});
Here, the reconnect behavior becomes part of the session configuration. reconnectInterval sets the initial reconnect interval to one second. maxReconnectInterval limits how long the interval can grow. backoffMultiplier: 2 increases the reconnect interval after unsuccessful attempts instead of retrying continuously at the same interval. Finally, maxRetries: 0 allows the session to keep retrying rather than stopping after a fixed number of reconnect attempts.
The practical difference in our application is that we can remove the custom post-connection recovery code from the broker. We no longer need connection.on("close", reconnect); or our custom reconnect() function. We also don’t need to manually call await subscribe(messageHandler); from a reconnect function. Our application subscribes normally, and the session owns the post-connection recovery behavior.
In our amqplib version, we wrote and maintained the post-connection recovery logic ourselves. In our amqp-client.js version, we configure that behavior and let the high-level session manage it. That means less connection lifecycle code in our broker and fewer recovery details for our application to manage.
Step 9: Replace shutdown
With amqplib, we close the channel and connection:
if (channel) await channel.close();
if (connection) await connection.close();
Replace that with:
if (session) {
await session.stop();
}
The rest of the application still calls await broker.close(); — no change is needed outside the broker.
Step 10: Select the client
Once both broker implementations expose the same functions, we can choose which one to use from the root:
const client = process.env.AMQP_CLIENT || "amqplib";
if (client === "amqp-client") {
module.exports = require("./amqp-client/broker");
} else if (client === "amqplib") {
module.exports = require("./amqplib/broker");
} else {
throw new Error(
`Unsupported AMQP_CLIENT "${client}". Use "amqplib" or "amqp-client".`
);
}
We then use our npm scripts to select the client. In package.json:
"produce:amqplib": "AMQP_CLIENT=amqplib node producer.js",
"worker:amqplib": "AMQP_CLIENT=amqplib node worker.js",
"produce:amqp-client": "AMQP_CLIENT=amqp-client node producer.js",
"worker:amqp-client": "AMQP_CLIENT=amqp-client node worker.js"
For example:
npm run worker:amqplib
runs the worker with the amqplib broker implementation. And:
npm run worker:amqp-client
runs the same worker with the amqp-client.js broker implementation.
Step 11: Test the migration
Let’s first confirm that the original implementation still works:
npm run worker:amqplib
In another terminal:
npm run produce:amqplib -- valid
Then switch to amqp-client.js:
npm run worker:amqp-client
And publish the same transaction:
npm run produce:amqp-client -- valid
We can test the failure paths as well:
npm run produce:amqp-client -- temporary-failure
npm run produce:amqp-client -- invalid
npm run produce:amqp-client -- review
The application behavior should remain the same. The final test is especially important for this migration: connection recovery. While the amqp-client.js worker is running, interrupt the connection to LavinMQ and then restore it. The worker should recover without us restarting the Node.js process manually.
This gives us a practical way to see the difference we discussed earlier: our amqplib implementation uses the recovery code we wrote, while our amqp-client.js implementation uses the recovery behavior configured on AMQPSession.
What actually changed?
The migration is mostly a small set of client-specific replacements:
amqplib |
amqp-client.js |
|---|---|
amqp.connect() |
AMQPSession.connect() |
connection.createChannel() |
Session-based high-level API |
channel.sendToQueue() |
queue.publish() |
channel.consume() |
queue.subscribe() |
channel.prefetch() |
prefetch subscription option |
message.content.toString() |
message.bodyString() |
channel.ack(message) |
message.ack() |
channel.publish() |
exchange.publish() |
| Our custom reconnect loop | Session recovery configuration |
| Close channel + connection | session.stop() |
- The producer didn’t change.
- The worker didn’t change.
- The transaction-processing logic didn’t change.
- The retry and dead-letter strategies didn’t change.
What changed is how our broker implementation communicates with LavinMQ and how it handles the connection lifecycle.
Summary
Migrating our Node.js LavinMQ application from amqplib to amqp-client.js doesn’t require us to rewrite the whole application. We replace the client-specific APIs for connecting, publishing, consuming, acknowledging messages, and shutting down.
The biggest behavioral difference in our implementation is connection recovery. With the amqplib broker, we listen for a closed connection, reconnect, recreate what we need, and subscribe the consumer again. With the amqp-client.js broker, we configure recovery on AMQPSession and let the high-level client manage that connection lifecycle for us.
The result is the same application running on top of a different LavinMQ client, with less recovery code for us to maintain.
Keep the application behavior. Replace the client-specific broker implementation.
Find the completed amqp-client/broker.js on GitHub.