Aller au contenu
English

Utiliser Redis et BullMQ

Configurez une file Redis partagée entre producteurs et workers de processus ou de machines différents.

Installez BullMQ à côté d’Outpost pour connecter les producteurs et les workers à la file Redis.

npm install bullmq
  • bullmqUn paquet optionnel, installé à côté d’Outpost.
  • Un serveur Redis autonomeAuto-hébergé ou managé, joignable par chaque producteur et chaque worker.
  • La politique noevictionExigée par la file, qui la vérifie à l’ouverture.

Utilisez une file Redis dédiée aux producteurs et aux workers qui partagent ces jobs. Configurez sa politique d’éviction avant de connecter BullMQ, afin que Redis ne supprime pas les clés de la file sous la pression mémoire.

Activez la persistance Redis (AOF ou snapshots RDB) si les jobs doivent survivre à un redémarrage de Redis.

Sur votre propre serveur, réglez la politique et conservez-la après redémarrage :

redis-cli CONFIG SET maxmemory-policy noeviction
redis-cli CONFIG REWRITE

Sur un service managé, modifiez-la dans la console : sur Redis Cloud, réglez Data eviction policy de la base sur No eviction ; sur Amazon ElastiCache, associez un groupe de paramètres personnalisé où maxmemory-policy vaut noeviction. Vérifiez ensuite ce que Redis indique :

redis-cli INFO memory | grep maxmemory_policy
# maxmemory_policy:noeviction

createBullMQTaskQueue() lit la même ligne INFO et refuse toute autre politique, ou un serveur qui la masque. Il ne modifie jamais la configuration du serveur.

Importez l’adaptateur depuis son propre sous-chemin. La file respecte le même contrat que la file SQLite : un worker s’y exécute sans changement (Files de jobs et workers).

import { runQueueWorker } from "@elie-laloum/outpost";
import { createBullMQTaskQueue } from "@elie-laloum/outpost/queues/bullmq";

const queue = await createBullMQTaskQueue({
  name: "code-reviews",
  connection: { host: "127.0.0.1", port: 6379 },
});
const stop = new AbortController();
process.once("SIGINT", () => stop.abort());
try {
  await runQueueWorker({
    queue,
    worker: `reviewer-${process.pid}`,
    signal: stop.signal,
    handlers: { review: (input) => ({ value: input }) },
  });
} finally {
  await queue.close();
}

Un producteur ouvre la même file et publie un job pour le traitement review :

import { createBullMQTaskQueue } from "@elie-laloum/outpost/queues/bullmq";

const queue = await createBullMQTaskQueue({
  name: "code-reviews",
  connection: { host: "127.0.0.1", port: 6379 },
});
try {
  await queue.enqueue({
    id: "review-42",
    handler: "review",
    input: { commit: "abc123" },
  });
} finally {
  await queue.close();
}

Référence API : BullMQTaskQueueOptions.

Chaque producteur et chaque worker doit utiliser les mêmes name, prefix, base Redis et serveur. Une seule différence lui donne, sans erreur, une file séparée et vide.

import { createBullMQTaskQueue } from "@elie-laloum/outpost/queues/bullmq";

const queue = await createBullMQTaskQueue({
  name: "code-reviews",
  prefix: "outpost",
  connection: {
    host: process.env.REDIS_HOST,
    port: 6379,
    db: 0,
    username: process.env.REDIS_USERNAME,
    password: process.env.REDIS_PASSWORD,
    tls: {},
  },
});
await queue.close();

Passez des réglages de connexion, pas un client Redis. Utilisez prefix plutôt que connection.keyPrefix, que l’adaptateur refuse. Réservez ce préfixe à Outpost : aucun autre consommateur BullMQ ni job de nettoyage ne doit toucher à ses clés.

  • ConnexionsUne pour ouvrir la file, puis une file et un worker BullMQ par traitement, ouverts au premier usage.
  • Recherche de baux expirésUn minuteur qui remet dans la file les jobs dont le bail a expiré.
  • Fermetureclose() refuse les nouveaux appels, attend ceux en cours, puis ferme chaque connexion.

Les jobs, baux et résultats restent dans Redis après close(), pour le processus suivant. Arrêtez le worker avant de fermer : annulez son signal et attendez runQueueWorker(), comme dans worker.ts.

La file enregistre chaque résultat dans son propre état Redis, puis marque le job BullMQ comme terminé. Ce résultat enregistré fait foi :

  • La finalisation échouecomplete() réussit quand même, get() renvoie le résultat et le job ne s’exécute plus jamais. L’erreur part vers onError.
  • enqueue() échoue en cours de routeRenvoyez la même requête avec le même id ; la file l’accepte et la publie.
  • Un worker planteSon bail expire et un autre worker reprend le job avec la même idempotencyKey.

Les réglages de connexion sont lus une seule fois, à l’ouverture de la file. Faites-les tourner en remplaçant les processus :

Glissez pour vous déplacer · Ctrl + molette pour zoomer
100 %
  • PréparerAvant tout redémarrage.
    1. Ajouter un second utilisateur ACLAvec les mêmes permissions que l’actuel. Redis
    (Étapes)
    • → Déployer : puis
  • DéployerAnciens et nouveaux processus partagent la file.
    1. Démarrer les nouveaux processusWorkers et producteurs avec le nouvel identifiant, les mêmes name et prefix. createBullMQTaskQueue()
    (Étapes)
    • → Retirer : puis
  • RetirerUne fois les nouveaux processus démarrés.
    1. Arrêter les anciens workersAnnuler leur signal, attendre runQueueWorker(), puis close(). runQueueWorker()
    2. RévoquerSupprimer l’ancien utilisateur ACL et déconnecter ses clients restants. Redis
    (Étapes)

Un identifiant révoqué pendant qu’un worker tourne encore lui fait perdre son bail. Un autre worker exécute alors le job de nouveau avec la même idempotencyKey : votre service d’effets doit dédupliquer (Files de jobs et workers).

  • Redis sans cluster : Redis Cluster n’est pas pris en charge. Le comportement lors d’une bascule de primaire, avec Sentinel ou un service managé, n’est pas garanti ; testez-le avant de vous y fier.
  • Effets au moins une fois : Les baux empêchent un worker périmé d’écrire un résultat, pas de répéter un effet externe. Dédupliquez avec idempotencyKey.
  • La durabilité est celle de Redis : Sans persistance, un redémarrage de Redis perd la file.

API : createBullMQTaskQueue · BullMQTaskQueueOptions · runQueueWorker.