🔄 Work Queues : Distribuez vos Tâches !
Imaginez que votre site web reçoit 1000 demandes de redimensionnement d'images par minute. Comment éviter que vos utilisateurs attendent 30 secondes pour chaque photo ? Les Work Queues sont la solution !
🎯 Le Problème à Résoudre
Avant les Work Queues
// ❌ Traitement synchrone = utilisateur qui attend
app.post('/upload', async (req, res) => {
const image = req.file;
// L'utilisateur attend 15 secondes...
await resizeImage(image, 'thumbnail');
await resizeImage(image, 'medium');
await resizeImage(image, 'large');
await optimizeImage(image);
res.json({ message: 'Image traitée!' });
});
Problème : L'utilisateur attend que tout soit fini avant de recevoir une réponse.
Avec les Work Queues
// ✅ Traitement asynchrone = réponse immédiate
app.post('/upload', async (req, res) => {
const image = req.file;
// Envoie la tâche dans une queue (instantané)
await sendTask('process_image', {
imageId: image.id,
operations: ['thumbnail', 'medium', 'large', 'optimize']
});
// Réponse immédiate !
res.json({ message: 'Traitement en cours!', imageId: image.id });
});
Avantage : L'utilisateur reçoit une réponse instantanée, le traitement se fait en arrière-plan.
🏗️ Comment ça Marche ?
Le Principe
- L'application envoie des tâches dans une queue
- Plusieurs workers traitent les tâches en parallèle
- Chaque tâche est traitée par un seul worker
- La charge est distribuée automatiquement
💻 Implémentation Pratique
Étape 1 : Créer le Producer (Envoie des Tâches)
Fichier : taskSender.js
const amqp = require('amqplib');
class TaskSender {
constructor() {
this.connection = null;
this.channel = null;
}
async connect() {
this.connection = await amqp.connect('amqp://localhost');
this.channel = await this.connection.createChannel();
// Créer la queue (survit aux redémarrages)
await this.channel.assertQueue('tasks', {
durable: true
});
}
async sendTask(taskType, data) {
const task = {
id: Date.now(),
type: taskType,
data: data,
createdAt: new Date()
};
// Envoyer la tâche (survit aux redémarrages)
await this.channel.sendToQueue('tasks',
Buffer.from(JSON.stringify(task)),
{ persistent: true }
);
console.log(`📤 Tâche envoyée: ${taskType} (ID: ${task.id})`);
return task.id;
}
}
module.exports = TaskSender;
Étape 2 : Utiliser dans votre Application
const express = require('express');
const TaskSender = require('./taskSender');
const app = express();
const taskSender = new TaskSender();
// Initialiser la connexion au démarrage
taskSender.connect();
app.post('/process-image', async (req, res) => {
const { imageUrl, userId } = req.body;
// Envoyer la tâche (très rapide)
const taskId = await taskSender.sendTask('process_image', {
imageUrl: imageUrl,
userId: userId,
operations: ['resize', 'optimize', 'watermark']
});
// Réponse immédiate
res.json({
message: 'Image en cours de traitement',
taskId: taskId
});
});
app.listen(3000, () => {
console.log('🚀 API lancée sur le port 3000');
});
Étape 3 : Créer le Worker (Traite les Tâches)
Fichier : worker.js
const amqp = require('amqplib');
class TaskWorker {
constructor(workerName) {
this.name = workerName;
this.connection = null;
this.channel = null;
}
async start() {
this.connection = await amqp.connect('amqp://localhost');
this.channel = await this.connection.createChannel();
// S'assurer que la queue existe
await this.channel.assertQueue('tasks', { durable: true });
// IMPORTANT: Une tâche à la fois par worker
this.channel.prefetch(1);
console.log(`🔄 ${this.name} démarré - En attente de tâches...`);
// Écouter les tâches
this.channel.consume('tasks', async (message) => {
if (message) {
await this.processTask(message);
}
}, { noAck: false });
}
async processTask(message) {
const task = JSON.parse(message.content.toString());
console.log(`🔨 ${this.name} traite: ${task.type} (ID: ${task.id})`);
try {
// Traiter selon le type de tâche
await this.handleTask(task);
// Tâche réussie
this.channel.ack(message);
console.log(`✅ ${this.name} terminé: ${task.type} (ID: ${task.id})`);
} catch (error) {
console.error(`❌ ${this.name} échoué:`, error.message);
// Remettre en queue pour retry
this.channel.nack(message, false, true);
}
}
async handleTask(task) {
switch (task.type) {
case 'process_image':
await this.processImage(task.data);
break;
case 'send_email':
await this.sendEmail(task.data);
break;
case 'generate_report':
await this.generateReport(task.data);
break;
default:
throw new Error(`Type de tâche inconnu: ${task.type}`);
}
}
async processImage(data) {
console.log(`📸 Traitement image: ${data.imageUrl}`);
// Simulation du traitement (2-5 secondes)
const processingTime = Math.random() * 3000 + 2000;
await new Promise(resolve => setTimeout(resolve, processingTime));
console.log(`📸 Image traitée pour utilisateur ${data.userId}`);
}
async sendEmail(data) {
console.log(`📧 Envoi email vers: ${data.email}`);
// Simulation envoi email
await new Promise(resolve => setTimeout(resolve, 1000));
console.log(`📧 Email envoyé: ${data.subject}`);
}
async generateReport(data) {
console.log(`📊 Génération rapport: ${data.reportType}`);
// Simulation génération rapport (plus long)
const processingTime = Math.random() * 5000 + 3000;
await new Promise(resolve => setTimeout(resolve, processingTime));
console.log(`📊 Rapport généré pour ${data.userId}`);
}
}
// Lancer le worker
const workerName = process.argv[2] || `Worker-${Math.random().toString(36).substring(7)}`;
const worker = new TaskWorker(workerName);
worker.start();
🚀 Test de votre Système
Lancer plusieurs workers
# Terminal 1 - Worker A
node worker.js "Worker-Images"
# Terminal 2 - Worker B
node worker.js "Worker-Emails"
# Terminal 3 - Worker C
node worker.js "Worker-Reports"
# Terminal 4 - API
node app.js
Tester l'API
# Envoyer plusieurs tâches
curl -X POST http://localhost:3000/process-image \
-H "Content-Type: application/json" \
-d '{"imageUrl": "photo1.jpg", "userId": 123}'
curl -X POST http://localhost:3000/process-image \
-H "Content-Type: application/json" \
-d '{"imageUrl": "photo2.jpg", "userId": 456}'
Vous verrez les tâches se distribuer automatiquement entre vos workers ! 🎉
🎛️ Configuration Avancée
Fair Dispatch (Distribution Équitable)
// ✅ Chaque worker reçoit une tâche à la fois
channel.prefetch(1);
// ❌ Un worker rapide pourrait recevoir toutes les tâches
// (pas de prefetch)
Messages avec Priorité
// Créer une queue avec priorités
await channel.assertQueue('priority_tasks', {
durable: true,
arguments: {
'x-max-priority': 10 // Priorité max = 10
}
});
// Envoyer avec priorité élevée
await channel.sendToQueue('priority_tasks',
Buffer.from(JSON.stringify(urgentTask)),
{
persistent: true,
priority: 9 // Priorité élevée
}
);
Gestion des Erreurs
async processTask(message) {
const task = JSON.parse(message.content.toString());
try {
await this.handleTask(task);
this.channel.ack(message); // ✅ Succès
} catch (error) {
console.error('Erreur:', error.message);
// Vérifier le nombre de tentatives
const retryCount = (message.properties.headers['x-retry-count'] || 0) + 1;
if (retryCount <= 3) {
// Remettre en queue avec compteur
await this.channel.sendToQueue('tasks', message.content, {
persistent: true,
headers: { 'x-retry-count': retryCount }
});
}
this.channel.ack(message); // Retirer de la queue
}
}
📊 Surveillance et Métriques
Monitoring Simple
class TaskWorker {
constructor(name) {
this.name = name;
this.stats = {
processed: 0,
failed: 0,
startTime: Date.now()
};
}
async processTask(message) {
const start = Date.now();
try {
await this.handleTask(task);
this.stats.processed++;
const duration = Date.now() - start;
console.log(`✅ Tâche terminée en ${duration}ms`);
} catch (error) {
this.stats.failed++;
console.error(`❌ Échec après ${Date.now() - start}ms`);
}
// Afficher les stats toutes les 10 tâches
if ((this.stats.processed + this.stats.failed) % 10 === 0) {
this.showStats();
}
}
showStats() {
const uptime = Date.now() - this.stats.startTime;
const rate = (this.stats.processed / (uptime / 1000)).toFixed(2);
console.log(`📊 ${this.name} - Traitées: ${this.stats.processed}, Échouées: ${this.stats.failed}, Taux: ${rate}/sec`);
}
}
🎯 Cas d'Usage Concrets
1. E-commerce : Traitement des Commandes
// Nouvelle commande
await sendTask('process_order', {
orderId: 12345,
items: [{ productId: 'ABC', quantity: 2 }],
customerId: 789,
actions: ['inventory_check', 'payment', 'shipping_label', 'email_confirmation']
});
2. Réseaux Sociaux : Traitement des Médias
// Upload vidéo
await sendTask('process_video', {
videoId: 'vid_123',
operations: ['extract_thumbnail', 'convert_formats', 'generate_subtitles'],
resolutions: ['480p', '720p', '1080p']
});
3. Analytics : Génération de Rapports
// Rapport mensuel
await sendTask('monthly_report', {
companyId: 456,
month: '2024-01',
includes: ['sales', 'traffic', 'conversions'],
format: 'pdf',
recipients: ['ceo@company.com', 'cto@company.com']
});
✅ Checklist des Bonnes Pratiques
Configuration
- [ ] Queues durables (
durable: true) - [ ] Messages persistants (
persistent: true) - [ ] Fair dispatch (
prefetch: 1) - [ ] Acknowledgments manuels (
noAck: false)
Gestion d'Erreurs
- [ ] Try/catch dans les workers
- [ ] Système de retry avec limite
- [ ] Dead Letter Queues pour les échecs
- [ ] Logs détaillés des erreurs
Performance
- [ ] Monitoring du nombre de tâches en attente
- [ ] Métriques de temps de traitement
- [ ] Alertes si les queues grandissent trop
- [ ] Possibilité d'ajouter des workers dynamiquement
Sécurité
- [ ] Validation des données de tâches
- [ ] Authentification RabbitMQ
- [ ] Isolation par Virtual Hosts
- [ ] Chiffrement des messages sensibles
🚀 Prochaines Étapes
Maintenant que vous maîtrisez les Work Queues, vous pouvez :
- Implémenter dans votre projet existant
- Ajouter des priorités aux tâches urgentes
- Créer des workers spécialisés par type de tâche
- Mettre en place le monitoring avec Grafana
- Explorer les patterns avancés (Pub/Sub, RPC)
💡 Astuce Pro : Commencez petit avec une seule queue et un seul type de tâche, puis étendez progressivement !
Félicitations ! 🎉 Vous pouvez maintenant gérer des millions de tâches asynchrones comme un pro. Vos utilisateurs ne remarqueront même plus les traitements longs !