Qu’est-ce qu’une file d’attente de tâches (job queue) ?
Une file d’attente de tâches est un tampon durable situé entre le code qui demande un travail et le code qui l’exécute. Le producteur écrit un petit enregistrement décrivant une unité de travail et revient immédiatement. Un worker lit cet enregistrement plus tard et s’occupe de la partie lente. La file d’attente est l’élément qui survit à un redémarrage, absorbe les pics de charge et permet de gérer les échecs.
Redis est utilisé de cette manière depuis plus d’une décennie. Ses listes ont fourni aux premières bibliothèques de files d’attente des primitives atomiques de LPUSH et BRPOP, et les “blocking pops” ont permis aux workers de rester en veille jusqu’à l’arrivée d’un travail au lieu de faire du polling en boucle. Les Streams ont ensuite ajouté les groupes de consommateurs, les accusés de réception et le replay. En s’appuyant sur ces primitives, BullMQ encapsule l’ensemble du cycle de vie — tâches différées, tentatives (retries), priorités, limitation de débit (rate limiting), planifications répétables et événements — dans une petite API TypeScript.
Si Redis gère déjà votre cache ou vos sessions, l’ajout d’une file d’attente est une étape simple. Cette commodité est la raison principale pour laquelle c’est le choix par défaut des équipes Node.js, et la raison majeure pour laquelle il est important de comprendre à la fois ce qu’il vous apporte et ses limites.
Une file d’attente, c’est trois rôles, pas un seul
Tout système de file d’attente, quel que soit le broker, repose sur les trois mêmes rôles, et c’est le fait de bien les distinguer qui rend le système maintenable.
Le producer est tout code qui ajoute un job. Il connaît le nom du job et la structure du payload, et rien d’autre. Il doit être rapide, car il s’exécute généralement au sein d’une requête, et son appel doit être idempotent (sans danger s’il est appelé deux fois).
app.post("/signup", async (req, res) => {
const user = await db.user.create({ data: req.body });
await emails.add("welcome", { userId: user.id, to: user.email });
res.status(201).json({ id: user.id });
});
La queue est l’état partagé au centre. Dans un déploiement Redis, il s’agit de l’ensemble des clés maintenues par BullMQ : la liste d’attente, l’ensemble des jobs différés, l’ensemble des jobs actifs et l’ensemble des jobs échoués. Elle persiste les jobs entre les redémarrages, les distribue de manière atomique et suit le nombre de tentatives.
Le worker est un processus de longue durée qui récupère les jobs et exécute les handlers. Il est séparé de l’API pour de bonnes raisons : il peut être déployé sur du matériel optimisé pour le CPU, être mis à l’échelle en fonction de la profondeur de la queue, et être redémarré sans perdre de requêtes. Un seul processus peut héberger des workers pour plusieurs queues, et une seule queue peut être servie par plusieurs processus worker. La queue est le seul état partagé, c’est pourquoi elle s’adapte si proprement à une mise à l’échelle horizontale.
Pourquoi Redis est un store de stockage courant
Trois propriétés de Redis en font une file d’attente naturelle.
Premièrement, les primitives existent déjà. Une liste avec LPUSH et BRPOP constitue une file d’attente, et BRPOP bloque la connexion plutôt que de consommer du CPU. C’est précisément pour cela que les premières bibliothèques de file d’attente pour Node ne faisaient que quelques centaines de lignes.
LPUSH queue:emails "welcome:42"
BRPOP queue:emails 30
Deuxièmement, les commandes sont atomiques. Déplacer un job de l’état d’attente à l’état actif, incrémenter son nombre de tentatives et définir son verrou peut tout se faire sans condition de concurrence (race condition), car Redis exécute les commandes une par une. C’est exactement ce dont une file d’attente a besoin pour éviter que deux workers ne récupèrent le même job.
Troisièmement, Redis est probablement déjà présent. C’est le cache et le store de session par défaut pour les services Node. Le réutiliser pour une file d’attente permet d’éviter d’ajouter une deuxième infrastructure, un deuxième jeu d’identifiants et un deuxième guide d’exploitation (runbook).
Les Streams poussent ce concept encore plus loin. Là où une liste ne peut être que poussée (push) et extraite (pop), un stream est un journal en ajout uniquement (append-only log) avec des groupes de consommateurs, des accusés de réception par message et un historique rejouable.
XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 STREAMS jobs:emails '>'
Le point faible reste la durabilité. Redis est avant tout un store en mémoire, donc un job acquitté mais pas encore écrit sur disque peut être perdu si l’instance s’arrête. Vous pouvez réduire cette fenêtre avec l’AOF et un replica, mais vous ne pourrez jamais rendre Redis aussi durable qu’une base de données avec write-ahead log. Considérez un job comme un travail récupérable, et non comme votre système d’enregistrement officiel.
Prise en main de BullMQ
BullMQ nécessite une connexion Redis autorisée à réessayer indéfiniment. Le comportement par défaut de ioredis abandonne une commande après quelques échecs, ce qui ne convient pas à un worker qui doit pouvoir surmonter une micro-coupure ; configurez donc maxRetriesPerRequest: null.
import { Queue, Worker } from "bullmq";
import IORedis from "ioredis";
const connection = new IORedis(process.env.REDIS_URL!, {
maxRetriesPerRequest: null,
});
const emails = new Queue("emails", { connection });
const worker = new Worker(
"emails",
async (job) => {
await sendEmail(job.data.to, job.data.template);
},
{ connection, concurrency: 10 },
);
Le Queue est l’identifiant du producteur. Le worker est le consommateur. Le nom "emails" correspond à la queue, et un worker ne voit que les jobs ajoutés à sa propre queue. Les queues distinctes constituent l’unité d’isolation : une queue d’importation lente ne peut pas retarder une queue de réinitialisation de mot de passe.
Un Queue est léger et peut être créé partout où vous avez besoin d’ajouter un job. Un Worker est un processus de longue durée et doit être lancé une seule fois par processus, et non par requête. Partagez l’objet de connexion entre toutes les queues et tous les workers du processus.
Ajouter des jobs avec des données et des options
add prend un nom, un payload et un objet d’options. Le nom route le job vers un handler ; le payload contient tout ce dont le worker a besoin.
await emails.add("welcome", { userId, to }, {
jobId: `welcome:${userId}`,
});
Le payload doit être un instantané de l’intention, et non un objet vivant. Si un utilisateur modifie son e-mail entre l’envoi dans la file d’attente et l’exécution, le job doit toujours être envoyé à l’adresse utilisée lors de sa création. Stockez les IDs et les quelques valeurs qui définissent le travail plutôt que la ligne complète de la base de données, car la file d’attente conserve chaque job en attente en mémoire.
L’objet d’options est là où réside la majeure partie de la valeur de BullMQ :
attempts— le nombre maximum de tentatives avant que le job ne soit considéré comme échoué.backoff— la stratégie de délai entre les tentatives, commeexponential.delay— ne pas exécuter avant ce nombre de millisecondes à partir de maintenant.priority— un nombre plus bas s’exécute en priorité lorsque des jobs sont en attente.jobId— un ID stable qui déduplique les envois du même travail logique.removeOnComplete— combien de jobs terminés conserver, outruepour les supprimer immédiatement.removeOnFail— s’il faut conserver les échecs pour inspection. Conservez-les.
Traitement des jobs et retour de résultats
Un handler reçoit le job, effectue le travail et peut retourner une valeur. La valeur de retour est stockée sur le job et peut être lue ultérieurement, ce qui transforme la file d’attente en un simple RPC asynchrone.
const worker = new Worker(
"reports",
async (job) => {
const pdf = await renderPdf(job.data.reportId);
return { url: pdf.url, bytes: pdf.bytes };
},
{ connection },
);
L’appelant peut ensuite interroger le job pour obtenir son résultat, ou écouter le flux d’événements pour savoir quand il est terminé.
const job = await reports.getJob(jobId);
if (await job?.isCompleted()) {
return job.returnvalue;
}
if (await job?.isFailed()) {
throw new Error(job.failedReason);
}
Retournez des valeurs de petite taille. Un résultat est stocké dans Redis comme n’importe quelle autre donnée ; retourner un PDF d’un mégaoctet encombrerait la file d’attente. Retournez plutôt une URL ou un id et laissez l’appelant récupérer les octets depuis un stockage d’objets.
Tentatives, backoff exponentiel et l’ensemble des échecs (failed set)
Les pannes transitoires sont normales. Une base de données bascule sur un serveur de secours, une API vous limite (rate-limiting), un conteneur est redéployé. Réessayer est la bonne réponse, mais réessayer immédiatement ne l’est pas.
await queue.add("sync", { accountId }, {
attempts: 5,
backoff: { type: "exponential", delay: 2_000 },
});
Cela produit des délais d’environ 2s, 4s, 8s et 16s, avec l’application du jitter propre à BullMQ pour éviter que de nombreux jobs ayant échoué simultanément ne soient relancés en même temps. Sans jitter, une flotte de workers touchés par la même panne recréerait un pic de charge dès le rétablissement du service.
Lorsqu’un job épuise ses tentatives, BullMQ ne le supprime pas. Il le déplace vers le failed set, accompagné du message d’erreur et de la trace de la pile (stack trace). Cet ensemble constitue la surface opérationnelle de la file d’attente : configurez des alertes lorsqu’il s’agrandit et mettez en place un mécanisme de rejeu (replay path).
const failed = await queue.getFailed(0, 20);
for (const job of failed) {
// Fix the underlying cause first, then requeue.
await job.retry();
}
Un job qui échoue systématiquement est un “poison message”. Le relancer indéfiniment est pire que de ne pas le relancer du tout, car il consomme un slot de worker à chaque tentative et prive les jobs sains de ressources. Limitez le nombre de tentatives et inspectez le failed set plutôt que de l’ignorer.
Jobs différés et répétitifs
Un job différé est programmé pour plus tard sans nécessiter de scheduler séparé. Un job répétitif s’exécute selon un pattern cron et appartient à la queue.
await queue.add("reminder", { userId }, { delay: 600_000 });
await queue.add(
"digest",
{ region: "eu" },
{
repeat: { pattern: "0 7 * * *", tz: "Europe/Berlin" },
jobId: "digest:eu",
},
);
L’jobId stable d’un job répétitif est crucial. Cela empêche le scheduler d’empiler une nouvelle copie si l’une est encore en cours d’exécution, et permet à chaque instance de l’application d’enregistrer le même planning sans créer de doublons. Privilégiez l’UTC ou un fuseau horaire explicite pour le pattern : un digest envoyé à 07:00 UTC n’est pas la même chose qu’un envoi à 07:00 heure locale, et cette différence se traduit par un ticket de support à chaque changement d’heure.
Les jobs différés sont stockés dans un sorted set indexé par leur heure d’exécution, ainsi, un job différé n’occupe pas de worker. Il reste dans Redis jusqu’à son échéance, ce qui rend les délais de plusieurs heures ou jours très peu coûteux en ressources.
Concurrence, limitation de débit et backpressure
La concurrence d’un worker correspond au nombre de jobs qu’il traite simultanément. L’augmenter permet d’accroître le débit jusqu’à ce que le worker manque de CPU, de connexions à la base de données ou de mémoire, moment où la situation se dégrade.
const worker = new Worker("sync", handler, {
connection,
concurrency: 10,
limiter: { max: 50, duration: 1_000 },
});
Le limiter limite le nombre de jobs que le worker démarre par fenêtre de temps. C’est la première ligne de défense pour éviter de saturer une API en aval. Si un fournisseur autorise 50 requêtes par seconde, un worker avec une concurrence de 200 vous exposera à un rate-limiting ; un limiteur réglé à 50 par seconde, non.
La backpressure est ce qui empêche la file d’attente elle-même de croître indéfiniment. Si les producteurs ajoutent des jobs plus rapidement que les workers ne peuvent les traiter, la file devient un backlog permanent et sa latence se compte en heures. Surveillez la profondeur de la file, mettez les producteurs en pause au-delà d’un certain seuil et scalez vos workers automatiquement. Une file d’attente qui ne fait que croître est une panne que personne n’a encore remarquée.
Progression et événements
Les tâches longues doivent rapporter leur progression afin qu’une interface utilisateur puisse afficher une barre de progression et qu’un opérateur puisse faire la différence entre une tâche lente et une tâche bloquée.
const worker = new Worker("imports", async (job) => {
const rows = await loadRows(job.data.fileId);
for (let i = 0; i < rows.length; i += 500) {
await insertChunk(rows.slice(i, i + 500));
await job.updateProgress(Math.round((i / rows.length) * 100));
}
return { rows: rows.length };
}, { connection });
Les événements sont le moyen par lequel le reste du système observe la file d’attente. Un worker émet completed, failed, progress et stalled pour les tâches qu’il exécute. QueueEvents écoute ce même flux depuis l’extérieur du worker, ce qui permet à un processus API de réagir à la fin d’une tâche sans être celui qui l’a exécutée.
worker.on("progress", (job, progress) => {
console.log(`job ${job.id} at ${progress}%`);
});
worker.on("completed", (job) => {
console.log(`job ${job.id} finished`);
});
Arrêt progressif (graceful shutdown) et jobs bloqués
Un worker interrompu en plein milieu d’un job laisse ce dernier dans un état ambigu. BullMQ gère cela avec un verrou (lock) : tant qu’un job est actif, le worker renouvelle un verrou dans Redis. Si le worker s’arrête et que le verrou expire, le job est marqué comme stalled (bloqué), replacé en attente et repris. C’est l’une des raisons pour lesquelles la livraison est de type “at-least-once”.
Comme un blocage entraîne une nouvelle exécution, gérez SIGTERM avec soin afin que les jobs en cours se terminent au lieu d’être interrompus.
process.on("SIGTERM", async () => {
await worker.close(); // stop accepting, wait for in-flight jobs
await connection.quit();
process.exit(0);
});
Accordez au déploiement un délai de grâce suffisamment long pour le job le plus lent, et limitez les timeouts des jobs pour qu’aucun ne puisse dépasser ce délai. Un handler pouvant s’exécuter pendant une heure devrait enregistrer sa progression (checkpoint) afin qu’une nouvelle exécution reprenne là où elle s’était arrêtée plutôt que de redémarrer depuis le début.
Utiliser Redis Streams directement
BullMQ est une couche d’abstraction basée sur les primitives de Redis, et il arrive que vous ayez besoin de ces primitives. Les Streams sont l’outil idéal lorsque plusieurs consommateurs doivent chacun voir chaque message, lorsque vous devez rejouer l’historique, ou lorsque vous souhaitez un accusé de réception sans passer par un framework de gestion de jobs.
XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 BLOCK 5000 STREAMS jobs:emails '>'
XACK jobs:emails workers 1760000000000-0
Le groupe de consommateurs suit les messages qui ont été livrés et ceux qui ont été acquittés. Un message livré mais jamais acquitté reste dans la liste d’attente (pending list) du groupe, afin qu’un consommateur ayant planté ne le perde pas. XAUTOCLAIM permet de réassigner les messages en attente d’un consommateur inactif à un consommateur actif.
Utilisez les streams directement lorsque les données constituent un journal d’événements et que les consommateurs sont des lecteurs indépendants. Utilisez BullMQ lorsque la donnée est une unité de travail avec une politique de retry, une priorité et un résultat. Réimplémenter les retries, les délais et un ensemble de jobs échoués (failed set) par-dessus les streams revient exactement à refaire le travail déjà accompli par BullMQ.
Idempotence : la livraison est de type « at-least-once »
Le point le plus important concernant les files d’attente est que la livraison est at-least-once (au moins une fois), et non exactement une fois. Un worker peut planter après avoir effectué le travail mais avant de l’avoir acquitté, et le job sera alors exécuté à nouveau. Par conception, un blocage entraîne une réexécution. Une livraison « exactly-once » à travers un réseau est concrètement impossible ; les files d’attente privilégient donc le « at-least-once » et vous transfèrent la responsabilité.
Votre handler doit être idempotent : l’exécuter deux fois doit produire le même état final que de l’exécuter une seule fois. La technique habituelle consiste à utiliser une clé de déduplication dérivée du travail, écrite de manière atomique avant l’effet de bord.
export async function handleCharge(job) {
const key = `charged:${job.data.orderId}`;
const first = await connection.set(key, "1", "NX", "EX", 86_400);
if (first === null) return { skipped: true, reason: "already_processed" };
await stripe.charges.create(
{ amount: job.data.amount, source: job.data.token },
{ idempotencyKey: job.data.orderId },
);
}
Remarquez que la clé provient de l’ID de la commande, et non de l’ID du job. Cela rend le handler sûr, même si le producteur place deux fois le même travail logique dans la file. Combinez la clé de déduplication avec la propre clé d’idempotence du fournisseur, car cette dernière ne protège que l’appel API, et non la logique environnante.
Surveiller la file d’attente
Une file d’attente est invisible à moins de la rendre visible. Quatre indicateurs couvrent la majorité de vos besoins.
- Profondeur de la file (Queue depth) — le nombre de jobs en attente. Une profondeur croissante signifie que les workers ne parviennent pas à suivre la cadence.
- Job en attente le plus ancien — la profondeur indique la quantité, l’âge indique la gravité. Dix mille jobs traités en une seconde ne posent aucun problème ; dix jobs qui attendent depuis une heure, si.
- Taux d’échec — le nombre d’échecs par minute, ventilé par nom de job. Un pic après un déploiement pointe directement vers la modification effectuée.
- Durée des jobs — un histogramme par nom de job. Un p95 en hausse signifie qu’une dépendance ralentit.
BullMQ expose ces compteurs directement, et un listener QueueEvents ou un exportateur peut les transmettre à votre système de métriques.
const counts = await queue.getJobCounts(
"wait", "active", "completed", "failed", "delayed",
);
console.log(counts);
Pour une vue visuelle, Bull Board ou le tableau de bord BullMQ montent une petite interface web sur les mêmes clés. Ajoutez un identifiant de corrélation (correlation id) à chaque payload et incluez-le dans les logs afin qu’un job puisse être tracé depuis la requête qui l’a créé jusqu’à chaque tentative de retry.
Garder des payloads légers et versionnés
Une file d’attente est une interface persistante entre deux déploiements. Un producteur exécutant la version 1 du code peut écrire un job que devra lire un worker exécutant la version 2. C’est le même problème de compatibilité que pour une API, et il est facile de l’ignorer jusqu’à ce qu’un déploiement ne casse un backlog.
Deux habitudes permettent de maintenir la compatibilité des payloads. Premièrement, ajoutez des champs plutôt que de les renommer ou de les supprimer, et attribuez une valeur par défaut cohérente aux nouveaux champs dans le handler. Un worker qui tolère l’absence de locale peut traiter des jobs mis en file d’attente avant que le champ n’existe. Deuxièmement, insérez une version dans le payload lorsque sa structure peut changer significativement, et utilisez un branchement conditionnel dans le handler.
await queue.add("import", { version: 2, importId, mapping });
Il y a une seconde raison de garder des payloads légers : Redis conserve chaque job en attente en mémoire. Un payload qui embarque une ligne entière de base de données se multiplie sur des milliers de jobs et rend la file d’attente coûteuse. Référencez les données par id et laissez le worker les récupérer. La seule exception est une valeur qui doit être figée au moment de la mise en file d’attente, comme le destinataire d’un e-mail ou le prix communiqué à un client ; celle-ci doit figurer dans le payload précisément parce qu’elle ne doit pas changer.
Composer des tâches avec les flows
Certains jobs sont en réalité composés de plusieurs étapes. Un rapport peut nécessiter la récupération de données, le rendu d’un PDF et son envoi par e-mail, et vous souhaiterez peut-être que chaque étape puisse être relancée indépendamment. Les flows de BullMQ permettent d’exprimer cela sous la forme d’un arbre de jobs parents et enfants.
import { FlowProducer } from "bullmq";
const flow = new FlowProducer({ connection });
await flow.add({
name: "report",
queueName: "reports",
data: { reportId },
children: [
{ name: "fetch", queueName: "reports", data: { reportId } },
{ name: "render", queueName: "reports", data: { reportId } },
],
});
Un job enfant s’exécute en premier, et le parent ne devient exécutable que lorsque tous ses enfants sont terminés. Cela vous permet de mettre en œuvre des mécanismes de fan-out et fan-in sans avoir à coordonner l’état dans votre propre base de données. Si un enfant échoue, le parent attend ou échoue selon les options du flow, ainsi la politique de retry reste attachée à l’étape qui a réellement posé problème. Gardez vos flows peu profonds et explicites ; un arbre complexe de jobs interdépendants est plus difficile à appréhender qu’un petit pipeline avec des étapes claires.
Choisir entre les files d’attente Redis, Kafka et RabbitMQ
Redis est le choix par défaut idéal lorsque l’unité de travail est un job : une tâche nommée avec une charge utile (payload), une politique de tentative (retry policy) et un résultat. C’est rapide, familier et déjà présent dans la plupart des stacks.
Kafka est le meilleur choix lorsque les données sont un log. Si plusieurs consommateurs indépendants doivent chacun lire chaque événement, si vous devez rejouer l’historique à partir d’un offset, ou si le débit se mesure en millions de messages par seconde, un log partitionné est plus adapté qu’une file d’attente de jobs.
RabbitMQ est le meilleur choix lorsque le routage est la partie complexe. Les exchanges et les binding keys permettent à un message d’être diffusé vers plusieurs files d’attente selon un modèle, avec des accusés de réception par message, des priorités et des dead-letter exchanges comme fonctionnalités natives. C’est plus lourd à exploiter que Redis, mais c’est un atout pour les équipes qui ont besoin de cette flexibilité.
Le conseil le plus honnête est de commencer avec les outils que vous utilisez déjà. Une file d’attente Redis est bien préférable à l’absence de file d’attente parce que vous attendiez d’évaluer différents brokers. Changez d’outil lorsqu’une limitation spécifique — durabilité, rejeu ou routage — devient réellement bloquante.
Bonnes pratiques
- Limitez vos gestionnaires de requêtes à une écriture et un enfilement, puis retournez
202 Accepted. - Partagez une seule connexion Redis par processus, avec
maxRetriesPerRequest: null. - Dérivez
jobIdà partir du travail pour que les enfilements en double soient fusionnés en un seul job. - Rendez chaque gestionnaire idempotent, car la livraison est de type “at-least-once”.
- Utilisez un backoff exponentiel avec jitter, et définissez une limite de tentatives.
- Conservez les échecs en configurant
removeOnFail: false, et limitez les jobs terminés avecremoveOnComplete. - Séparez les files d’attente par charge de travail afin que les jobs lents n’empêchent pas l’exécution des jobs urgents.
- Limitez la concurrence et ajoutez un limiteur pour toute dépendance soumise à un rate-limit.
- Rapportez la progression des jobs longs pour différencier les jobs lents de ceux qui sont bloqués.
- Arrêtez les workers sur
SIGTERMet prévoyez un délai de grâce correspondant lors des déploiements. - Suivez la profondeur de la file, l’âge du job le plus ancien, le taux d’échec et la durée sur un tableau de bord.
Erreurs courantes
- Supposer qu’un job s’exécute exactement une seule fois, et ainsi facturer un client deux fois.
- Utiliser un id de job aléatoire et accumuler des plannings répétitifs en double.
- Retenter indéfiniment un “poison message” et saturer la queue.
- Exécuter des tâches lourdes dans le processus de l’API et appeler cela une queue.
- Régler la concurrence si haut que la base de données atteint sa limite de connexions.
- Retourner un objet énorme comme résultat de job et saturer Redis.
- Oublier
maxRetriesPerRequest: nullet perdre des jobs lors d’une reconnexion. - Ignorer l’ensemble des jobs échoués sur tous les tableaux de bord jusqu’au premier incident.
- Tuer les workers avec
SIGKILLet forcer chaque job en cours à stagner. - Placer une ligne entière de la base de données dans le payload au lieu d’un id.
Pour aller plus loin
Le guide Redis détaille le stockage utilisé par la file d’attente : les listes, les streams, les TTL, la persistance et le modèle de commandes mono-thread qui permet des revendications atomiques. Pour approfondir la pratique consistant à sortir le travail du flux de requête, consultez Batch Processing. Lorsque le routage et les accusés de réception par message deviennent complexes, passez à RabbitMQ, et si vous avez besoin d’un journal rejouable plutôt que d’une file de tâches, le guide Kafka est l’étape suivante.