🔍 RabbitMQ vs Kafka vs Redis

Introduction

Choisir la bonne solution de messaging est crucial pour l'architecture d'un système distribué. RabbitMQ, Apache Kafka et Redis Pub/Sub ont chacun leurs forces et cas d'usage optimaux.

📊 Vue d'Ensemble Comparative

Rendu du diagramme en cours...

🐰 RabbitMQ

Caractéristiques Principales

# Architecture RabbitMQ
"""
Strengths:
- Protocol AMQP complet
- Routing flexible (exchanges)
- Acknowledgments et transactions
- Interface de management UI
- Plugins écosystème riche
"""

# Exemple typique RabbitMQ
import pika

class RabbitMQExample:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
    
    def setup_complex_routing(self):
        # Multiple exchange types
        self.channel.exchange_declare('direct_ex', 'direct')
        self.channel.exchange_declare('topic_ex', 'topic')
        self.channel.exchange_declare('fanout_ex', 'fanout')
        
        # Dead Letter Exchange
        self.channel.exchange_declare('dlx', 'direct')
        
        # Queue avec TTL et DLX
        self.channel.queue_declare(
            queue='task_queue',
            durable=True,
            arguments={
                'x-message-ttl': 60000,
                'x-dead-letter-exchange': 'dlx',
                'x-max-length': 10000
            }
        )

Forces de RabbitMQ

  1. Routing Complexe
// Topic exchange avec wildcards
channel.publish('events', 'order.eu.created', orderData);
channel.publish('events', 'order.us.shipped', shipmentData);

// Consumers avec patterns
subscribe('order.*.created');  // Tous les created
subscribe('order.eu.*');       // Tous les EU
subscribe('#.shipped');         // Tous les shipped
  1. Garanties de Livraison
# Publisher confirms
channel.confirm_delivery()

# Consumer acknowledgments
def callback(ch, method, properties, body):
    try:
        process_message(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
  1. Management et Monitoring
# API REST pour monitoring
curl -u guest:guest http://localhost:15672/api/queues

Cas d'Usage Optimaux

  • ✅ Communication inter-services
  • ✅ Task queues / Job processing
  • ✅ RPC patterns
  • ✅ Routing complexe
  • ✅ Guaranteed delivery

🚀 Apache Kafka

Caractéristiques Principales

// Architecture Kafka
/*
Strengths:
- Event streaming platform
- Log immutable distribué
- Haute performance (millions msg/sec)
- Replay de messages
- Écosystème complet (Streams, Connect)
*/

// Producer Kafka
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

Producer<String, String> producer = new KafkaProducer<>(props);

// Envoi avec partitioning
ProducerRecord<String, String> record = new ProducerRecord<>(
    "events",           // topic
    event.getUserId(),  // key for partitioning
    event.toJson()      // value
);

producer.send(record);

Forces de Kafka

  1. Event Streaming & Replay
from kafka import KafkaConsumer

# Consumer avec replay depuis le début
consumer = KafkaConsumer(
    'events',
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='earliest',  # Replay depuis le début
    group_id='analytics-group'
)

# Ou depuis un offset spécifique
consumer.seek(TopicPartition('events', 0), 1000)
  1. Haute Performance
# Benchmarks typiques
performance:
  throughput: 1000000 msg/sec
  latency: < 10ms p99
  storage: compression jusqu'à 10x
  retention: jours/semaines/mois
  1. Stream Processing
// Kafka Streams
KStream<String, Order> orders = builder.stream("orders");

KTable<String, Long> orderCounts = orders
    .filter((key, order) -> order.getStatus().equals("COMPLETED"))
    .groupByKey()
    .count();

orderCounts.toStream().to("order-counts");

Cas d'Usage Optimaux

  • ✅ Event sourcing
  • ✅ Log aggregation
  • ✅ Stream processing
  • ✅ Data pipeline
  • ✅ Analytics temps réel

🔴 Redis Pub/Sub

Caractéristiques Principales

import redis

# Architecture Redis Pub/Sub
"""
Strengths:
- Ultra faible latence
- Simple à implémenter
- In-memory (très rapide)
- Structures de données riches
- Léger
"""

class RedisExample:
    def __init__(self):
        self.redis_client = redis.Redis(
            host='localhost',
            port=6379,
            decode_responses=True
        )
    
    def publisher(self):
        # Publication simple
        self.redis_client.publish('notifications', 'Hello subscribers!')
        
        # Avec Redis Streams (plus robuste)
        self.redis_client.xadd(
            'events:stream',
            {'user': 'john', 'action': 'login'}
        )
    
    def subscriber(self):
        pubsub = self.redis_client.pubsub()
        pubsub.subscribe('notifications')
        
        for message in pubsub.listen():
            if message['type'] == 'message':
                print(f"Received: {message['data']}")

Forces de Redis

  1. Performance Ultra-Rapide
// Benchmarks Redis
const redis = require('redis');
const client = redis.createClient();

// Pub/Sub: < 1ms latency
await client.publish('channel', 'message');

// Redis Streams pour plus de garanties
await client.xAdd('mystream', '*', {
  sensor: 'temp01',
  value: '23.5'
});
  1. Structures de Données
# Au-delà du pub/sub
redis_client.lpush('queue:tasks', task_json)  # List as queue
redis_client.zadd('leaderboard', {user: score})  # Sorted set
redis_client.hset('session:123', mapping=session_data)  # Hash
redis_client.setex('cache:user:456', 3600, user_json)  # TTL

Cas d'Usage Optimaux

  • ✅ Real-time notifications
  • ✅ Chat applications
  • ✅ Live updates
  • ✅ Cache avec invalidation
  • ✅ Session management

📈 Comparaison Détaillée

Performance

| Critère | RabbitMQ | Kafka | Redis | |---------|----------|-------|--------| | Throughput | 20-50K msg/s | 1M+ msg/s | 100K+ msg/s | | Latence | 1-10ms | 10-100ms | <1ms | | Persistance | Oui (disque) | Oui (log) | Optionnel | | Scalabilité | Vertical + Horizontal | Horizontal infini | Vertical principalement |

Fonctionnalités

Rendu du diagramme en cours...

Garanties de Livraison

# RabbitMQ - At least once avec ack
channel.basic_consume(
    queue='my_queue',
    on_message_callback=callback,
    auto_ack=False  # Manual ack
)

# Kafka - Configurable
producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    acks='all',  # Wait for all replicas
    enable_idempotence=True  # Exactly once
)

# Redis Pub/Sub - At most once (fire & forget)
redis_client.publish('channel', 'message')  # No guarantees

# Redis Streams - At least once
redis_client.xreadgroup(
    'mygroup',
    'consumer1',
    {'mystream': '>'},
    block=0
)

🎯 Matrice de Décision

Choisir RabbitMQ quand :

use_rabbitmq_when:
  - routing_complexe: true
  - different_exchange_types: true
  - rpc_pattern: true
  - priority_queues: true
  - dead_letter_handling: true
  - moderate_throughput: "10K-50K msg/s"
  - guaranteed_delivery: critical
  - management_ui: required

Choisir Kafka quand :

use_kafka_when:
  - event_sourcing: true
  - log_aggregation: true
  - stream_processing: true
  - high_throughput: "> 100K msg/s"
  - message_replay: required
  - long_retention: "days to months"
  - ordered_processing: true
  - horizontal_scaling: critical

Choisir Redis quand :

use_redis_when:
  - ultra_low_latency: "< 1ms"
  - simple_pub_sub: true
  - ephemeral_messages: true
  - caching_layer: true
  - real_time_updates: true
  - small_messages: true
  - in_memory_speed: required

🏗️ Architectures Hybrides

Combiner les Solutions

// Architecture hybride commune
class HybridMessaging {
  constructor() {
    // Redis pour cache et notifications temps réel
    this.redis = new Redis();
    
    // RabbitMQ pour processing asynchrone
    this.rabbit = amqp.connect('amqp://localhost');
    
    // Kafka pour event sourcing
    this.kafka = new Kafka({
      clientId: 'hybrid-app',
      brokers: ['localhost:9092']
    });
  }
  
  async processOrder(order) {
    // 1. Cache in Redis for fast access
    await this.redis.set(`order:${order.id}`, JSON.stringify(order), 'EX', 3600);
    
    // 2. Send to RabbitMQ for processing
    await this.rabbit.publish('orders', 'order.created', order);
    
    // 3. Log to Kafka for event sourcing
    await this.kafka.producer.send({
      topic: 'order-events',
      messages: [{
        key: order.id,
        value: JSON.stringify({
          type: 'OrderCreated',
          timestamp: Date.now(),
          data: order
        })
      }]
    });
    
    // 4. Notify via Redis Pub/Sub
    await this.redis.publish('order-notifications', `New order: ${order.id}`);
  }
}

📊 Cas d'Usage Réels

E-Commerce Platform

# RabbitMQ pour workflow de commande
def process_order_workflow(order):
    # Complex routing entre services
    channel.publish('orders.created', order)
    # → Payment Service
    # → Inventory Service
    # → Shipping Service

# Kafka pour analytics
def log_user_activity(event):
    # Stream d'événements pour ML/Analytics
    producer.send('user-activity', event)

# Redis pour panier temps réel
def update_cart(user_id, item):
    redis.hset(f'cart:{user_id}', item.id, item.to_json())
    redis.publish(f'cart-updates:{user_id}', 'item_added')

IoT Platform

// Redis pour données temps réel des capteurs
redis.xadd('sensors:temperature', '*', {
  sensor_id: 'temp_01',
  value: 23.5,
  timestamp: Date.now()
});

// Kafka pour stockage long terme
kafka.send({
  topic: 'sensor-data',
  messages: [{ key: sensorId, value: sensorData }]
});

// RabbitMQ pour alertes
if (temperature > threshold) {
  rabbitmq.publish('alerts', 'temperature.high', {
    sensor: sensorId,
    value: temperature
  });
}

🔄 Migration Entre Solutions

De RabbitMQ vers Kafka

# Bridge RabbitMQ → Kafka
class RabbitToKafkaBridge:
    def bridge_messages(self):
        def rabbit_callback(ch, method, properties, body):
            # Forward to Kafka
            self.kafka_producer.send(
                'migrated-topic',
                key=method.routing_key,
                value=body
            )
            ch.basic_ack(delivery_tag=method.delivery_tag)
        
        self.rabbit_channel.basic_consume(
            queue='source-queue',
            on_message_callback=rabbit_callback
        )

✅ Recommandations

  1. Start Simple : Redis Pub/Sub pour prototypes
  2. Scale Smart : RabbitMQ pour la plupart des cas
  3. Scale Big : Kafka pour Big Data/Streaming
  4. Combine : Architecture hybride pour le meilleur des mondes

Le choix dépend de vos besoins spécifiques en termes de performance, garanties, et complexité acceptable.

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours