🔄 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

Rendu du diagramme en cours...
  1. L'application envoie des tâches dans une queue
  2. Plusieurs workers traitent les tâches en parallèle
  3. Chaque tâche est traitée par un seul worker
  4. 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 :

  1. Implémenter dans votre projet existant
  2. Ajouter des priorités aux tâches urgentes
  3. Créer des workers spécialisés par type de tâche
  4. Mettre en place le monitoring avec Grafana
  5. 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 !

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours