Qu’est-ce que RabbitMQ ?
RabbitMQ est un message broker open-source qui utilise AMQP (Advanced Message Queuing Protocol). Les producteurs y publient des messages, les consommateurs les reçoivent, et le broker est responsable de la conservation de chaque message jusqu’à ce que celui-ci soit acquitté.
La phrase qui décrit le mieux sa conception est : un broker intelligent avec des consommateurs simples. RabbitMQ gère les exchanges, les règles de routage, les acquittements, les tentatives de réessai (retries) et les dead letters. Un consommateur est un petit programme qui lit un message, effectue une action et signale que le travail est terminé. Cette division du travail explique pourquoi RabbitMQ s’intègre parfaitement dans les systèmes où le routage et les garanties de livraison sont plus importants que le débit brut.
Il a été écrit en Erlang en 2007, c’est pourquoi il est exceptionnellement performant pour gérer de nombreuses connexions simultanées et pourquoi il jouit d’une solide réputation de stabilité. Il implémente un protocole complet plutôt qu’une API de file d’attente minimale, et ce protocole est à la fois la source de sa puissance et de sa courbe d’apprentissage.
Le modèle AMQP
L’idée fondamentale de RabbitMQ est qu’un producteur ne publie jamais directement dans une queue. Il publie vers un exchange, et c’est l’exchange qui décide quelles queues recevront une copie.
producer -> exchange --binding--> queue -> consumer
\--binding--> queue -> consumer
Le modèle repose sur cinq éléments :
- Un producteur ouvre un channel et appelle
basic.publishavec un nom d’exchange, une routing key et un corps de message. - Un exchange reçoit chaque message publié et l’achemine. Il ne stocke rien, à moins d’être lié à une queue.
- Un binding est une règle reliant un exchange à une queue. Il peut contenir une routing key ou un pattern, et un seul exchange peut être lié à plusieurs queues.
- Une queue stocke les messages dans l’ordre jusqu’à ce qu’un consommateur les récupère.
- Un consommateur s’abonne à une queue, reçoit les messages et en accuse réception.
C’est tout l’intérêt de cette indirection. Un producteur qui publie order.created.eu ne sait pas si un service, cinq services ou personne n’écoute. On ajoute de nouveaux consommateurs en créant un nouveau binding vers une nouvelle queue, sans jamais modifier le producteur.
Tout le travail s’effectue sur un channel, qui est une connexion virtuelle légère multiplexée sur une seule connexion TCP. Les channels ne sont pas thread-safe ; le pattern habituel consiste donc à utiliser un channel par tâche ou par consommateur.
import amqp from "amqplib";
const conn = await amqp.connect(process.env.AMQP_URL!);
const ch = await conn.createChannel();
Types d’exchange et routage
Le type d’un exchange détermine la manière dont il associe une routing key à ses bindings. Il en existe quatre, chacun ayant un usage précis.
Un direct exchange route les messages vers les queues dont la binding key est exactement identique à la routing key. Utilisez-le pour envoyer un message à une queue spécifique, ou à un petit ensemble de queues partageant toutes la même clé.
await ch.assertExchange("logs", "direct", { durable: true });
await ch.bindQueue("logs.errors", "logs", "error");
await ch.publish("logs", "error", body);
Un fanout exchange ignore complètement la routing key et copie le message vers toutes les queues liées. C’est le pub/sub dans sa forme la plus pure : une publication, plusieurs consommateurs indépendants.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.bindQueue("search-indexer", "events", "");
await ch.bindQueue("email-notifier", "events", "");
Un topic exchange compare la routing key à un pattern. Les clés sont composées de mots séparés par des points. * correspond exactement à un mot et # correspond à zéro ou plusieurs mots, ce qui fait des topic exchanges les plus flexibles des quatre.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.bindQueue("eu-orders", "orders", "order.*.eu");
await ch.bindQueue("all-orders", "orders", "order.#");
await ch.publish("orders", "order.created.eu", body);
Un headers exchange ignore la routing key et effectue l’association sur la base des attributs du header du message. C’est rarement le choix optimal, car les topic exchanges sont plus faciles à lire et à appréhender, mais cela s’avère utile lorsque le routage dépend de plusieurs attributs indépendants.
Si vous ne devez retenir qu’une seule règle, c’est celle-ci : choisissez le type d’exchange en fonction de la question que vous vous posez. « Quelle queue précise ? » $\rightarrow$ direct. « Tout le monde ? » $\rightarrow$ fanout. « Quelle famille d’événements ? » $\rightarrow$ topic.
Acknowledgements, nack et prefetch
La garantie de livraison de RabbitMQ repose sur les acquittements (acknowledgements). Lorsqu’un consommateur reçoit un message, le broker le marque comme unacked mais le conserve. Le message n’est supprimé que lorsque le consommateur appelle ack. Si la connexion est interrompue avant cela, le broker redélivre le message.
await ch.consume("orders.created", async (msg) => {
if (!msg) return;
try {
await handleOrder(JSON.parse(msg.content.toString()));
ch.ack(msg);
} catch (err) {
ch.nack(msg, false, false); // requeue: false, dead-letter instead
}
});
nack (ou l’ancienne méthode reject) accepte un flag requeue. Le requeueing replace le message en tête de file, ce qui est approprié pour une erreur transitoire, mais problématique pour un “poison message” qui échouera systématiquement. La solution habituelle consiste alors à l’envoyer vers un dead-letter exchange.
Le prefetch, configuré via basic.qos, limite le nombre de messages non acquittés qu’un consommateur peut détenir simultanément. Sans cela, le broker pousse les messages aussi vite que possible et un consommateur lent peut se retrouver avec des milliers de messages en mémoire.
await ch.prefetch(20);
Réglez le prefetch sur un petit multiple de votre concurrence réelle. S’il est trop élevé, un seul consommateur peut monopoliser la file ; s’il est trop bas, le consommateur restera inactif en attendant le prochain aller-retour.
Durabilité et garanties de livraison
Un message survit au redémarrage d’un broker uniquement si trois conditions sont réunies, et il est facile d’en oublier une.
- La queue doit être déclarée
durable: true. - Le message doit être publié avec
persistent: true. - L’exchange doit également être
durable: true.
await ch.assertExchange("orders", "topic", { durable: true });
await ch.assertQueue("orders.created", { durable: true });
ch.publish("orders", "order.created.eu", body, { persistent: true });
La durabilité concerne la survie après un redémarrage, et non la garantie de livraison. Pour cela, activez les publisher confirms. Sans confirms, publish fonctionne en mode “fire-and-forget” : si le broker s’arrête avant d’avoir écrit le message, le producteur n’en sera jamais informé.
const ch = await conn.createConfirmChannel();
ch.publish("orders", "order.created.eu", body, { persistent: true }, (err) => {
if (err) console.error("broker did not confirm", err);
else console.log("message is safely queued");
});
Même avec les confirms et des queues durables, la livraison est de type at-least-once (au moins une fois). Un consommateur peut planter après avoir effectué le travail mais avant d’avoir envoyé l’ack, et le broker redistribuera alors le message. Le “exactly-once” à travers un réseau est pratiquement impossible ; RabbitMQ rend donc ce compromis explicite et vous demande de rendre vos handlers idempotents.
Dead-letter exchanges et files d’attente de retry
Un dead-letter exchange (DLX) est un exchange ordinaire qui reçoit les messages qu’une file d’attente rejette, dont le délai d’expiration est atteint ou qui sont abandonnés. C’est le mécanisme qui permet de mettre en place des files d’attente de retry et la gestion des “poison messages”.
Un message est envoyé vers le DLX lorsque l’un des événements suivants se produit :
- Le consommateur envoie un nack ou le rejette avec
requeue: false. - Son TTL (Time To Live), défini par message ou par file d’attente, expire.
- La file d’attente dépasse sa limite de longueur et abandonne le message le plus ancien.
Le DLX se configure sur la file d’attente qui contient le travail, et non sur le consommateur.
await ch.assertQueue("orders.created", {
durable: true,
deadLetterExchange: "orders.dlx",
deadLetterRoutingKey: "failed",
});
Le pattern classique de retry différé utilise une seconde file d’attente avec un TTL et son propre DLX pointant vers l’exchange de travail. Un message en échec est envoyé vers la file d’attente de retry, y reste pendant la durée du TTL, expire, puis est renvoyé via le DLX pour être traité à nouveau. Cela permet d’effectuer un retry avec un délai sans avoir à gérer de timer dans votre code.
await ch.assertExchange("orders.dlx", "direct", { durable: true });
await ch.assertQueue("orders.retry", {
durable: true,
messageTtl: 30_000,
deadLetterExchange: "orders",
deadLetterRoutingKey: "order.created.retry",
});
await ch.bindQueue("orders.retry", "orders.dlx", "retry");
Un message qui échoue systématiquement entrera dans une boucle infinie. Suivez le nombre de tentatives dans les headers du message et, après avoir atteint une limite, routez-le vers une file d’attente failed permanente que personne ne consomme automatiquement. Considérez cette file d’attente comme une surface opérationnelle : configurez une alerte lorsqu’elle s’accumule et prévoyez un mécanisme de rejeu (replay).
TTL des messages et limites de file d’attente
Le Time-to-live (TTL) contrôle la durée pendant laquelle un message peut attendre. Vous pouvez le configurer par file d’attente, par message, ou les deux.
await ch.assertQueue("verification", {
durable: true,
messageTtl: 600_000, // 10 minutes for every message
maxLength: 10_000, // keep at most 10k messages
overflow: "reject-publish", // backpressure instead of dropping
});
Le TTL par message est défini lors de la publication et est souvent utilisé pour des valeurs qui varient d’un message à l’autre.
ch.publish("orders", key, body, {
expiration: "30000", // milliseconds, as a string
});
Les limites de longueur de file d’attente servent de backpressure. maxLength limite la mémoire, et overflow décide de ce qui se passe lorsque la limite est atteinte : drop-head rejette silencieusement le message le plus ancien, tandis que reject-publish refuse les nouvelles publications pour que le producteur ressente la pression. Pour une file d’attente qui ne doit perdre aucune tâche, reject-publish est le choix par défaut le plus sûr et sert de signal d’alerte.
Files d’attente de travail (work queues) vs pub/sub
Le même broker gère deux modes de communication très différents, et les confondre peut entraîner des bugs.
Une work queue distribue chaque message à un seul consommateur. Plusieurs workers se disputent les messages d’une même file d’attente, et le broker répartit les livraisons via un round-robin. En ajoutant un prefetch et des acquittements (acknowledgements), vous obtenez une file d’attente de tâches scalable et tolérante aux pannes.
await ch.assertQueue("jobs.thumbnails", { durable: true });
await ch.prefetch(5);
Le pub/sub livre chaque message à tous les consommateurs intéressés. Chaque abonné possède sa propre file d’attente liée au même exchange ; ainsi, un abonné lent ou hors ligne ne “vole” jamais de message aux autres.
await ch.assertExchange("events", "fanout", { durable: true });
await ch.assertQueue("events.billing", { durable: true });
await ch.bindQueue("events.billing", "events", "");
await ch.assertQueue("events.analytics", { durable: true });
await ch.bindQueue("events.analytics", "events", "");
La règle d’or : une seule file d’attente partagée par plusieurs consommateurs est une work queue ; une file d’attente par consommateur est du pub/sub. Un topic exchange vous permet de combiner les deux, avec certains consommateurs partageant une file d’attente et d’autres possédant la leur.
Le pattern RPC
RabbitMQ peut également gérer le mode requête/réponse. Le client publie une requête avec une file d’attente replyTo et un correlationId, et le serveur publie la réponse dans cette file d’attente avec le même identifiant.
import { randomUUID } from "node:crypto";
const { queue } = await ch.assertQueue("", { exclusive: true });
const correlationId = randomUUID();
ch.consume(queue, (msg) => {
if (msg?.properties.correlationId === correlationId) {
console.log("reply", JSON.parse(msg.content.toString()));
}
}, { noAck: true });
ch.publish("rpc.inventory", "check", Buffer.from("{}"), {
replyTo: queue,
correlationId,
});
Le assertQueue("") crée une file d’attente temporaire, exclusive et s’auto-supprimant, dédiée uniquement à ce client. Le RPC via une file d’attente est utile lorsque l’appelant a réellement besoin d’une réponse, mais cela réintroduit un couplage synchrone et un problème de timeout. Pour la plupart des systèmes, un événement suivi d’un événement de rappel (callback) est plus simple à exploiter que le RPC, et un simple appel HTTP est encore plus facile lorsque la dépendance est disponible.
Clustering et quorum queues
Un nœud unique constitue un point de défaillance unique (single point of failure), c’est pourquoi RabbitMQ en production s’exécute sous forme de cluster. Les queues peuvent être répliquées sur plusieurs nœuds, et les clients se reconnectent à un autre nœud lorsqu’un premier tombe en panne.
La stratégie de réplication moderne est la quorum queue, basée sur l’algorithme de consensus Raft. Une quorum queue possède un leader et des followers, et une écriture n’est confirmée qu’une fois qu’une majorité les a reçue. Cela la rend fiable face aux partitions réseau qui causaient des pertes de données avec les mirrored queues classiques, et c’est aujourd’hui le choix par défaut pour les queues durables.
await ch.assertQueue("orders.created", {
durable: true,
arguments: { "x-queue-type": "quorum" },
});
Les quorum queues privilégient un petit nombre impair de réplicas, généralement trois ou cinq. Elles sont plus lourdes que les queues classiques ; utilisez-les donc pour les données que vous ne pouvez pas perdre et conservez les queues transitoires en mode classique. Pour les topologies multi-régions, la federation et le plugin shovel permettent de déplacer des messages entre des brokers indépendants plutôt que d’étendre un seul cluster sur une liaison réseau lente.
L’interface de gestion et le monitoring
Chaque nœud RabbitMQ est livré avec un plugin de gestion qui fournit une interface web et une API HTTP. Elle permet de visualiser les exchanges, les queues, les bindings, les connexions et les channels, et vous permet de publier un message de test ou d’en rejouer un depuis une queue.
rabbitmq-plugins enable rabbitmq_management
curl -u guest:guest http://localhost:15672/api/queues/%2F/orders.created
Quatre indicateurs sont essentiels sur un tableau de bord.
- Queue depth — messages prêts. Une profondeur croissante signifie que les consumers ne parviennent pas à suivre la cadence.
- Unacked count — messages livrés mais non acquittés. Un chiffre qui ne fait que grimper indique que les consumers sont bloqués.
- Publish and deliver rates — la forme du trafic et la capacité des consumers à maintenir le rythme.
- Redelivery rate — une augmentation du nombre de redélivrances signale des crashs ou des échecs répétés.
RabbitMQ émet également des métriques Prometheus, ainsi la queue depth et l’utilisation des consumers doivent figurer sur le même tableau de bord que le reste de vos services. Configurez vos alertes sur la profondeur de la queue et l’âge du message le plus ancien, plutôt que sur des échecs individuels, qui sont attendus.
Connexions, canaux et récupération
Une connexion est coûteuse ; un canal est léger. Ouvrez une seule connexion par processus, puis créez un canal par producteur, par consommateur ou par unité de travail. Un canal n’est pas thread-safe : le partager entre plusieurs gestionnaires concurrents provoque l’entrelacement des trames et des erreurs confuses.
import amqp from "amqplib";
let conn: amqp.Connection;
let ch: amqp.Channel;
async function connect() {
conn = await amqp.connect(process.env.AMQP_URL!);
conn.on("error", (err) => console.error("connection error", err));
conn.on("close", () => setTimeout(connect, 5_000));
ch = await conn.createChannel();
await ch.assertExchange("orders", "topic", { durable: true });
}
Gérez la reconnexion de manière délibérée. amqplib ne gère pas la reconnexion pour vous, vous devez donc écouter l’événement close et reconstruire la connexion, le canal ainsi que chaque consommateur. Comme une connexion interrompue laisse des messages en cours non acquittés (unacked), le broker les redistribue, ce qui est le comportement souhaité. C’est un autre cas où la livraison “at-least-once” intervient : une reconnexion peut rejouer du travail, les gestionnaires doivent donc être idempotents.
Propriétés et priorités des messages
Chaque message publié peut transporter des propriétés en plus de son corps. Elles voyagent avec le message et sont visibles pour les consommateurs, ce qui en fait un endroit léger pour stocker des métadonnées.
ch.publish("orders", "order.created.eu", body, {
persistent: true,
contentType: "application/json",
messageId: order.id,
correlationId: traceId,
timestamp: Math.floor(Date.now() / 1000),
type: "order.created",
headers: { "x-retry-count": 0, "x-source": "checkout" },
});
contentType et messageId aident les consommateurs et l’outillage ; correlationId lie un message à une trace ; headers sont l’endroit où conserver les métadonnées de l’application, comme le nombre de tentatives de réessai (retry count). Ne placez pas de données volumineuses dans les headers, car le broker doit les indexer et les afficher.
RabbitMQ prend également en charge les priority queues (files d’attente prioritaires), où les messages ayant une priorité plus élevée sont livrés en premier. Déclarez la queue avec une priorité maximale et publiez avec une valeur priority.
await ch.assertQueue("jobs", {
durable: true,
maxPriority: 10,
});
ch.publish("", "jobs", body, { priority: 9 });
Les priorités ne réorganisent que les messages qui sont déjà en attente. Si les consommateurs maintiennent la queue vide, la priorité n’a aucun effet, et un message à haute priorité qui arrive après un message à basse priorité attendra toujours derrière lui. Utilisez la priorité pour une distinction métier réelle, et non comme un mécanisme d’ordonnancement général.
Choisir entre RabbitMQ et Kafka
On compare souvent les deux, mais ils répondent à des problématiques différentes.
RabbitMQ est un broker intelligent conçu pour le routage et les files d’attente de travail (work queues). Les messages sont des tâches qui sont supprimées une fois acquittées. Les exchanges effectuent le routage par pattern, les acquittements se font par message, et la gestion des dead-lettering est intégrée. Il excelle lorsqu’un message doit atteindre un ensemble spécifique de files d’attente, ou lorsqu’une unité de travail doit être retentée puis mise de côté en cas d’échec.
Kafka est un log durable. Les messages sont ajoutés à un log ordonné et partitionné, et sont conservés pendant une durée configurée, quel que soit le lecteur. Plusieurs groupes de consommateurs peuvent lire le même topic indépendamment, et tout consommateur peut revenir à un offset antérieur. Il excelle pour les débits très élevés et pour le replay d’événements.
Une heuristique utile : si le fait qu’un message soit consommé puis supprimé vous pose problème, vous avez besoin d’un log. Si vous voulez vous assurer qu’une tâche a été correctement routée et terminée, vous avez besoin d’un broker. De nombreux systèmes utilisent les deux : Kafka pour le flux d’événements et RabbitMQ pour le travail.
Bonnes pratiques
- Déclarez les exchanges, les queues et les bindings de manière idempotente au démarrage, et utilisez
durable: truepour tout élément critique. - Publiez des messages persistants vers des queues durables, et utilisez les publisher confirms pour les tâches qui ne doivent pas disparaître.
- Accusez toujours réception manuellement, une fois que l’effet de bord a réussi, et jamais avant.
- Définissez une limite de prefetch pour éviter qu’un seul consumer ne mette en tampon l’intégralité de la queue.
- Attribuez un dead-letter exchange à chaque queue durable, et mettez en place un chemin de rejeu (replay path).
- Implémentez des tentatives de rejeu différées avec une queue de retry TTL plutôt que d’utiliser un sleep dans le handler.
- Rendez vos consumers idempotents, car la livraison est de type at-least-once.
- Limitez la longueur des queues et choisissez
reject-publishlorsque la perte de messages est inacceptable. - Utilisez des quorum queues pour les données durables et des classic queues pour les données transitoires.
- Fermez les channels et les connexions lors de
SIGTERM, et arrêtez la consommation avant le drainage. - Surveillez les indicateurs ready, unacked, le taux de redelivery et l’âge du message le plus ancien.
Erreurs courantes
- Publier vers une queue au lieu d’un exchange, puis se demander pourquoi les bindings ne fonctionnent pas.
- Utiliser l’auto-ack (
noAck: true) et perdre des messages dès qu’un handler plante. - Oublier
persistent: true, et perdre ainsi des messages lors du redémarrage d’un broker. - Déclarer une queue comme durable tout en publiant des messages non persistants.
- Régler le prefetch trop haut, ce qui permet à un seul consumer d’affamer tous les autres.
- Re-queuer indéfiniment un “poison message” avec
nack(msg, false, true). - Créer un dead-letter exchange sans consumer et ne jamais aller le consulter.
- Utiliser un exchange fanout là où un exchange topic était nécessaire, ou inversement.
- Faire tourner un broker unique en production et prétendre que l’installation est hautement disponible.
- Mélanger les sémantiques de work-queue et de pub/sub sur une même queue.
- Bloquer l’event loop dans un consumer, ce qui interrompt les heartbeats et déclenche une déconnexion intempestive.
Et après ?
Si vous recherchez le même concept de travail fiable mais avec moins d’infrastructure, le guide sur les Redis Queues détaille l’utilisation de BullMQ sur Redis. Pour approfondir la discipline consistant à sortir les tâches du flux de requête, consultez la section Batch Processing. Lorsque vos événements constituent un journal durable rejoué par plusieurs consommateurs plutôt que des tâches routées puis supprimées, le guide Kafka est l’étape suivante logique, et les Node.js basics couvrent le runtime sur lequel s’exécute chaque consommateur amqplib.