Modern Backend Engineering
Chapitre 13
Chapitre 13 — Message Queues
Chapitre 13 — Message Queues
Cours — Message Queues
1. Introduction aux Message Queues
Pourquoi les files d'attente ?
- Découplage : producteur et consommateur ne se connaissent pas
- Buffer : lissage des pics de charge
- Async : traitement différé des tâches longues
- Scalabilité : multiples consommateurs
- Résilience : messages persistés en cas de panne
Patterns de messagerie
| Pattern | Description | Cas d'usage |
|---|---|---|
| Point-to-Point | 1 producteur → 1 consommateur | Task queue |
| Pub/Sub | 1 producteur → N consommateurs | Notifications |
| Request/Reply | Message avec réponse | RPC |
| Competing Consumers | N consommateurs sur 1 queue | Load balancing |
2. RabbitMQ (AMQP)
Concepts AMQP
Publisher → Exchange → Binding → Queue → Consumer
Types d'exchanges
| Type | Routage |
|---|---|
| Direct | routing_key exacte |
| Topic | routing_key avec pattern (topic.#, topic.*) |
| Fanout | broadcast à toutes les queues |
| Headers | basé sur les en-têtes |
Exemple Node.js
import amqp from 'amqplib';
async function setup() {
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
// Declare exchange
await channel.assertExchange('orders', 'topic', { durable: true });
// Declare queue
await channel.assertQueue('order.created', { durable: true });
// Bind queue to exchange
await channel.bindQueue('order.created', 'orders', 'order.created.*');
// Publish
channel.publish('orders', 'order.created.new', Buffer.from(JSON.stringify({
orderId: 123,
userId: 456,
amount: 99.99
})), {
persistent: true,
contentType: 'application/json'
});
// Consume
channel.consume('order.created', (msg) => {
if (msg) {
const data = JSON.parse(msg.content.toString());
console.log('Processing order:', data.orderId);
channel.ack(msg);
}
});
}
Dead Letter Queue (DLQ)
await channel.assertQueue('orders.dlq', { durable: true });
await channel.assertQueue('orders.main', {
durable: true,
deadLetterExchange: '',
deadLetterRoutingKey: 'orders.dlq',
maxLength: 1000,
messageTtl: 60000, // 1 minute
});
3. Kafka
Architecture
Topic (partition 0) ──▶ Consumer Group A
Topic (partition 1) ──▶ Consumer Group A
Topic (partition 2) ──▶ Consumer Group B
Concepts clés
- Topic : catégorie de messages
- Partition : division d'un topic (ordre garanti dans une partition)
- Offset : position du message dans la partition
- Consumer Group : groupe de consommateurs (load balancing)
- Broker : serveur Kafka
- Replication : copies des partitions (HA)
Producer
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
clientId: 'order-service',
brokers: ['localhost:9092'],
});
const producer = kafka.producer({
allowAutoTopicCreation: true,
transactionTimeout: 30000,
});
await producer.connect();
await producer.send({
topic: 'orders',
messages: [
{
key: 'order-123',
value: JSON.stringify({ orderId: 123, userId: 456 }),
headers: { 'event-type': 'order.created' },
},
],
});
await producer.disconnect();
Consumer
const consumer = kafka.consumer({ groupId: 'order-processor' });
await consumer.connect();
await consumer.subscribe({ topic: 'orders', fromBeginning: true });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
console.log({
topic,
partition,
offset: message.offset,
key: message.key.toString(),
value: message.value.toString(),
});
// Traitement
await processOrder(JSON.parse(message.value.toString()));
// Auto-commit (ou manuel avec enableAutoCommit: false)
},
autoCommitInterval: 5000, // commit toutes les 5s
});
Kafka Streams (traitement temps réel)
const { StreamsConfig } = require('kafkajs');
// Stream processing avec kafka.js ou utilisez kafka-streams
// Exemple conceptuel :
// 1. Lire depuis un topic source
// 2. Transformer les données
// 3. Écrire dans un topic destination
4. Redis Pub/Sub
import Redis from 'ioredis';
const publisher = new Redis();
const subscriber = new Redis();
// Subscriber
subscriber.subscribe('notifications', (err, count) => {
console.log(`Subscribed to ${count} channels`);
});
subscriber.on('message', (channel, message) => {
console.log(`Received ${message} from ${channel}`);
});
// Publisher
setInterval(() => {
publisher.publish('notifications', JSON.stringify({
type: 'alert',
message: 'Server CPU > 80%',
timestamp: Date.now(),
}));
}, 5000);
Redis Streams (persistant)
// Add to stream
await publisher.xadd('mystream', '*', 'key', 'value', 'temperature', 25);
// Read from stream
const results = await subscriber.xread(
'BLOCK', 5000,
'STREAMS', 'mystream', '$'
);
5. Architecture Event-Driven
Event-Driven Microservices
[Order Service] ──▶ order.created ──▶ [Payment Service]
│
├──▶ [Inventory Service]
│
└──▶ [Notification Service]
Avantages
- Découplage fort entre services
- Évolutivité indépendante
- Résilience (un service down n'impacte pas les autres)
- Traçabilité (event sourcing)
Inconvénients
- Complexité (eventual consistency)
- Debug plus difficile
- Nécessite de la monitoring
6. Dead Letter Queues (DLQ)
Quand utiliser une DLQ ?
- Message mal formaté
- Échec de traitement après N tentatives
- Timeout de traitement
- Message en dehors des limites (taille, TTL)
async function consumeWithRetry(channel, queue, dlq, maxRetries = 3) {
await channel.assertQueue(queue, {
deadLetterExchange: '',
deadLetterRoutingKey: dlq,
});
channel.consume(queue, async (msg) => {
try {
await processMessage(msg);
channel.ack(msg);
} catch (error) {
const retryCount = (msg.properties.headers['x-retry-count'] || 0) + 1;
if (retryCount <= maxRetries) {
// Retry with delay
channel.publish(
'',
queue,
msg.content,
{
headers: { 'x-retry-count': retryCount, 'x-delay': retryCount * 5000 },
persistent: true,
}
);
channel.ack(msg);
} else {
// Send to DLQ
channel.reject(msg, false); // false = ne pas requeue
}
}
});
}
7. Message Ordering et Partitioning
Garantir l'ordre
- Kafka : ordre garanti dans une partition
- RabbitMQ : ordre garanti dans une queue (avec 1 consommateur)
Partition key
// Kafka : même clé = même partition = ordre garanti
await producer.send({
topic: 'orders',
messages: [{ key: `user-${userId}`, value: payload }],
});
Problèmes d'ordre
- Plusieurs consommateurs sur une queue = ordre non garanti
- Solution : router par clé de partition
8. Idempotent Consumers
Principe
Un message peut être traité plusieurs fois → le résultat doit être le même.
Comment implémenter l'idempotence
async function processEvent(event) {
// 1. Vérifier si déjà traité
const alreadyProcessed = await redis.exists(`processed:${event.id}`);
if (alreadyProcessed) {
return; // Déjà traité
}
// 2. Marquer comme en cours
await redis.set(`processing:${event.id}`, Date.now(), 'EX', 30);
// 3. Traiter (transactionnel)
await db.transaction(async (tx) => {
// Insérer le résultat
await tx.query('INSERT INTO events (id, data) VALUES ($1, $2)', [event.id, event.data]);
// Marquer comme traité
await redis.set(`processed:${event.id}`, '1', 'EX', 86400);
});
// 4. Nettoyer
await redis.del(`processing:${event.id}`);
}
At-Least-Once vs Exactly-Once
| Guarantee | Description | Coût |
|---|---|---|
| At-most-once | Peut perdre des messages | Faible |
| At-least-once | Peut dupliquer | Moyen |
| Exactly-once | Pas de perte, pas de duplication | Élevé |
9. Comparaison des Message Brokers
| Critère | RabbitMQ | Kafka | Redis Pub/Sub |
|---|---|---|---|
| Modèle | Queue + Exchange | Log (partition) | Pub/Sub |
| Persistance | Oui | Oui | Non (ou Streams) |
| Ordre | 1 queue | 1 partition | Non garanti |
| Throughput | 10k msg/s | 1M+ msg/s | 100k msg/s |
| Rétention | Ack-based | Time-based | En mémoire |
| Cas d'usage | Tasks, RPC | Events, logs, streams | Real-time, cache |
| Complexité | Moyenne | Élevée | Faible |
10. Patterns Avancés
Event Sourcing
Stockage de l'état comme séquence d'événements.
- Avantage : audit trail complet, rejouer les events
- Inconvénient : complexité, stockage
CQRS (Command Query Responsibility Segregation)
Séparation des commandes (écritures) et requêtes (lectures).
- Commande : validation, mise à jour, publication d'event
- Requête : lecture optimisée (matérialisation)
Saga Pattern
Pour les transactions distribuées.
- Choreography : chaque service publie/reagit aux events
- Orchestration : un orchestrateur coordonne
Outbox Pattern
Pour garantir la fiabilité des événements.
1. Écrire dans la BDD + outbox (même transaction)
2. Processus lit l'outbox et publie les events
3. Supprime de l'outbox après publication