# Fichier: python_cheats/cheatsheets/kafka.txt
# Introduction Apache Kafka - Pour Grands Débutants
# Comprendre Kafka de A à Z : Architecture, Fonctionnement, Pratique

"""
[OK] INTRODUCTION : POURQUOI CE GUIDE ?

# Tu débutes avec les systèmes de messagerie distribués et Apache Kafka ?
# Ce guide est fait pour TOI !

# Objectif :
# Te montrer EXACTEMENT comment Kafka fonctionne "sous le capot"
# Comprendre ce qui se passe quand tu publies/consommes un message
# Découvrir TOUS les concepts et leur rôle précis
# Maîtriser l'architecture distribuée
# Aucun concept ne sera laissé dans le flou !

# Plan du guide :
# 1. Qu'est-ce que Kafka et pourquoi existe-t-il ?
# 2. Architecture fondamentale (Brokers, Topics, Partitions)
# 3. Installation et configuration détaillée
# 4. Producers : Écrire des messages
# 5. Consumers : Lire des messages
# 6. Consumer Groups et parallélisme
# 7. Kafka sous le capot (Log, Segments, Index)
# 8. Réplication et haute disponibilité
# 9. Kafka Streams (traitement en temps réel)
# 10. Kafka Connect (intégrations)
# 11. Schema Registry et sérialisation
# 12. Exactly-once semantics et transactions
# 13. Performance, tuning et monitoring
# 14. Cas d'usage réels et patterns
# 15. Exemples pratiques complets
"""


# [OK] PARTIE 1 : QU'EST-CE QUE KAFKA ET POURQUOI ?

"""
┌────────────────────────────────────────────────────────────────────────┐
│                    APACHE KAFKA : UNE RÉVOLUTION                       │
└────────────────────────────────────────────────────────────────────────┘

CONTEXTE HISTORIQUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AVANT KAFKA (années 2000) :
-> RabbitMQ, ActiveMQ : Message brokers traditionnels
-> Bases de données pour le stockage
-> ETL batch pour le traitement
-> Architecture point-à-point ou pub/sub simple

PROBLÈMES AVEC L'APPROCHE TRADITIONNELLE :
[X] Difficile de gérer des millions d'événements par seconde
[X] Couplage fort entre systèmes
[X] Pas de replay des messages (une fois consommé = perdu)
[X] Scaling vertical limité
[X] Latence élevée pour le traitement en temps réel
[X] Complexité de l'intégration de multiples sources

NAISSANCE DE KAFKA (2011 chez LinkedIn) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Créé par Jay Kreps, Neha Narkhede, Jun Rao
-> Open-sourcé en 2011
-> Apache Top-Level Project en 2012
-> Écrit en Scala/Java
-> Version actuelle : 3.6+ (décembre 2024)

Kafka = "log distribué" + "système de messagerie" + "plateforme de streaming"


QU'EST-CE QUE KAFKA ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION SIMPLE :
Kafka est une plateforme de streaming d'événements distribuée qui permet de :
1. PUBLIER et SOUSCRIRE à des flux de messages (comme un système de messagerie)
2. STOCKER ces flux de manière durable et distribuée (comme une base de données)
3. TRAITER ces flux en temps réel (comme un moteur de traitement de flux)


ANALOGIE : JOURNAL INTERNE D'UNE ENTREPRISE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Imagine une grande entreprise avec un JOURNAL INTERNE :

┌────────────────────────────────────────────────────────────────────┐
│ [NEWSPAPER] JOURNAL INTERNE DE L'ENTREPRISE                                 │
│                                                                    │
│ Chaque département publie ses événements dans le journal :         │
│                                                                    │
│ [NOTE] Service RH      -> "Nouveau employé embauché : Jean Dupont"     │
│ [ARGENT] Comptabilité    -> "Facture payée : #12345, 5000€"              │
│ [PACKAGE] Logistique      -> "Colis expédié : #98765, destination Dakar"  │
│ [SHOPPING_TROLLEY] Ventes          -> "Commande reçue : #54321, client XYZ"        │
│                                                                    │
│ Caractéristiques du journal :                                     │
│ [OK] Immuable : une fois écrit, on ne peut pas modifier             │
│ [OK] Ordonné : les événements sont dans l'ordre chronologique       │
│ [OK] Persistant : conservé pendant X jours/semaines/mois            │
│ [OK] Accessible : tous les départements peuvent le lire             │
│ [OK] Répliqué : plusieurs copies pour sécurité                      │
│                                                                    │
│ Chaque département peut :                                         │
│ -> Publier ses propres événements (PRODUCER)                       │
│ -> Lire les événements qui l'intéressent (CONSUMER)                │
│ -> Commencer à lire depuis n'importe quel point dans le temps      │
│ -> Relire les anciens événements si nécessaire                     │
└────────────────────────────────────────────────────────────────────┘

KAFKA = Ce journal interne, mais à l'échelle de millions d'événements/sec !


SYSTÈME DE MESSAGERIE TRADITIONNEL vs KAFKA :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MESSAGE QUEUE TRADITIONNELLE (RabbitMQ, ActiveMQ) :
┌────────────────────────────────────────────────────────────────────┐
│ [POSTBOX] FILE DE MESSAGES (QUEUE)                                        │
│                                                                    │
│ Producer -> [Msg1] [Msg2] [Msg3] -> Consumer                        │
│                                                                    │
│ Fonctionnement :                                                  │
│ 1. Producer envoie un message dans la queue                       │
│ 2. Consumer lit le message                                        │
│ 3. Message SUPPRIMÉ de la queue [X]                                │
│ 4. Impossible de relire le message                                │
│                                                                    │
│ Limites :                                                         │
│ [X] Message consommé = perdu pour toujours                         │
│ [X] Difficile de scaler (vertical surtout)                         │
│ [X] Pas de replay                                                  │
│ [X] Couplage temporel (consumer doit être up)                      │
└────────────────────────────────────────────────────────────────────┘


KAFKA (LOG DISTRIBUÉ) :
┌────────────────────────────────────────────────────────────────────┐
│ [LISTE] LOG DISTRIBUÉ (IMMUTABLE, APPEND-ONLY)                         │
│                                                                    │
│ Producer -> [Msg1][Msg2][Msg3][Msg4][Msg5]... -> Consumer A         │
│                   ^                                ^               │
│                   └────────────────────────────────┘ Consumer B    │
│                                                                    │
│ Fonctionnement :                                                  │
│ 1. Producer APPEND un message au log                              │
│ 2. Message stocké de manière PERSISTANTE (disque)                 │
│ 3. Consumer A lit le message à son offset                         │
│ 4. Consumer B lit le MÊME message à son offset                    │
│ 5. Messages GARDÉS pendant X jours (rétention)                    │
│ 6. Possibilité de REPLAY à tout moment                            │
│                                                                    │
│ Avantages :                                                       │
│ [OK] Messages persistés et réutilisables                            │
│ [OK] Scaling horizontal (partitions)                                │
│ [OK] Replay infini (dans période de rétention)                      │
│ [OK] Découplage total (consumers indépendants)                      │
│ [OK] Haut débit (millions de messages/sec)                          │
└────────────────────────────────────────────────────────────────────┘


POURQUOI KAFKA ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. DÉCOUPLAGE DES SYSTÈMES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AVANT KAFKA (couplage direct) :
┌──────────┐ ────-> ┌──────────┐
│  App A   │       │  App B   │
└──────────┘       └──────────┘
     │
     └─────────-> ┌──────────┐
                 │  App C   │
                 └──────────┘

-> App A doit connaître B et C
-> Si B tombe, A doit gérer l'erreur
-> Complexité = N * (N-1) connexions

AVEC KAFKA (découplage) :
┌──────────┐         ┌────────────────┐         ┌──────────┐
│  App A   │ ──────-> │     KAFKA      │ ──────-> │  App B   │
└──────────┘         │   (Topics)     │         └──────────┘
                     └────────────────┘
                            │
                            └─────-> ┌──────────┐
                                    │  App C   │
                                    └──────────┘

-> App A publie dans Kafka
-> B et C lisent de manière indépendante
-> Complexité = N connexions (linéaire)


2. SCALABILITÉ HORIZONTALE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Message Queue traditionnelle :
-> Scaling vertical (serveur plus puissant)
-> Limite physique atteinte rapidement

Kafka :
-> Scaling horizontal (ajouter des brokers)
-> Partitioning natif (parallélisme)
-> Millions de messages/sec possibles


3. DURABILITÉ ET RÉPLICATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

-> Messages persistés sur disque
-> Réplication sur plusieurs brokers
-> Pas de perte de données (si config correcte)
-> Rétention configurable (heures, jours, années)


4. PERFORMANCE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

-> Écritures séquentielles sur disque (très rapide)
-> Zero-copy (kernel -> network, pas de passage en userspace)
-> Batch processing
-> Compression native
-> Latence : quelques millisecondes


5. REPLAY ET REPROCESSING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Cas d'usage :
-> Bug dans le consumer -> Rembobiner et retraiter
-> Nouveau consumer -> Lire tous les événements passés
-> Analyse rétrospective -> Rejouer les données historiques
-> Machine learning -> Réentraîner sur données historiques


QUI UTILISE KAFKA ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

-> LinkedIn (créateur, 7+ trillions de messages/jour)
-> Netflix (traçage, logging, analytics)
-> Uber (tracking temps réel des courses)
-> Airbnb (événements utilisateurs)
-> Spotify (streaming de logs et métriques)
-> Twitter (analytics en temps réel)
-> PayPal (transactions et fraude)
-> Adidas (e-commerce events)
-> Goldman Sachs (trading)


CONCEPTS CLÉS KAFKA :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MESSAGE (ou RECORD ou EVENT) :
-> Unité de base de données dans Kafka
-> Composé de : Key, Value, Timestamp, Headers

TOPIC :
-> Catégorie de messages (comme un dossier)
-> Exemple : "orders", "user-clicks", "payments"

PARTITION :
-> Subdivision d'un topic pour le parallélisme
-> Chaque partition = log ordonné et immuable

PRODUCER :
-> Application qui PUBLIE des messages dans un topic

CONSUMER :
-> Application qui LIT des messages depuis un topic

BROKER :
-> Serveur Kafka qui stocke les messages

CLUSTER :
-> Ensemble de brokers Kafka

OFFSET :
-> Identifiant unique de position d'un message dans une partition


KAFKA vs AUTRES TECHNOLOGIES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────────┬────────────┬────────────┬────────────┬────────────┐
│  Caractéristique │  Kafka     │ RabbitMQ   │  Redis     │  Pulsar    │
├──────────────────┼────────────┼────────────┼────────────┼────────────┤
│ Type             │ Log distri │ Message    │ Key-Value  │ Log distri │
│                  │ bué        │ broker     │ store      │ bué        │
│ Persistance      │ Disque     │ Optionnelle│ RAM        │ Disque     │
│ Ordre garanti    │ Par        │ Par queue  │ Non        │ Par        │
│                  │ partition  │            │            │ partition  │
│ Débit            │ Très élevé │ Moyen      │ Très élevé │ Très élevé │
│ Latence          │ ~2-5ms     │ ~1ms       │ <1ms       │ ~5-10ms    │
│ Replay           │ Oui [OK]    │ Non [X]     │ Non [X]     │ Oui [OK]     │
│ Scalabilité      │ Horizontal │ Limité     │ Horizontal │ Horizontal │
│ Cas d'usage      │ Streaming, │ Task queue │ Cache,     │ Multi-     │
│                  │ Event      │            │ Pub/Sub    │ tenancy    │
│ Complexité       │ Moyenne    │ Faible     │ Faible     │ Élevée     │
│ Écosystème       │ Très riche │ Mature     │ Mature     │ Croissant  │
└──────────────────┴────────────┴────────────┴────────────┴────────────┘


QUAND UTILISER KAFKA ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] OUI SI :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Besoin de traiter des millions d'événements/sec
-> Event sourcing (stocker l'historique des événements)
-> CDC (Change Data Capture) - synchroniser les bases
-> Log aggregation (centraliser les logs)
-> Streaming analytics (analyse en temps réel)
-> Microservices communication
-> IoT (millions de capteurs)
-> Activity tracking (clicks, vues, actions)
-> Metrics et monitoring
-> Besoin de replay des données

[X] NON SI :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Simple task queue (RabbitMQ suffira)
-> Très faible latence critique (<1ms) (utiliser Redis)
-> Petite échelle (<1000 messages/sec)
-> Pas besoin de persistance
-> Pas d'équipe avec compétences distributed systems
-> Request-response pattern (utiliser REST/gRPC)


ARCHITECTURE GLOBALE SIMPLIFIÉE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

                        KAFKA CLUSTER
    ┌──────────────────────────────────────────────────┐
    │                                                  │
    │   ┌─────────┐  ┌─────────┐  ┌─────────┐          │
    │   │ Broker 1│  │ Broker 2│  │ Broker 3│          │
    │   └─────────┘  └─────────┘  └─────────┘          │
    │        │            │            │               │
    │        └────────────┴────────────┘               │
    │                    │                             │
    │              ┌─────────────┐                     │
    │              │  ZooKeeper  │                     │
    │              │   ou KRaft  │                     │
    │              └─────────────┘                     │
    └──────────────────────────────────────────────────┘
           [BLACK_UP-POINTING_TRIANGLE]                            │
           │                            │
    ┌──────────────┐            ┌──────────────┐
    │  PRODUCERS   │            │  CONSUMERS   │
    │              │            │              │
    │ ┌─────────┐  │            │ ┌─────────┐  │
    │ │ App A   │  │            │ │ App X   │  │
    │ └─────────┘  │            │ └─────────┘  │
    │ ┌─────────┐  │            │ ┌─────────┐  │
    │ │ App B   │  │            │ │ App Y   │  │
    │ └─────────┘  │            │ └─────────┘  │
    └──────────────┘            └──────────────┘

FLUX BASIQUE :
1. PRODUCERS envoient des messages aux BROKERS
2. BROKERS stockent les messages dans des TOPICS/PARTITIONS
3. CONSUMERS lisent les messages depuis les BROKERS
4. ZooKeeper (ou KRaft) coordonne le cluster
"""


# [OK] PARTIE 2 : ARCHITECTURE FONDAMENTALE

"""
┌────────────────────────────────────────────────────────────────────────┐
│              ARCHITECTURE KAFKA - CONCEPTS DÉTAILLÉS                   │
└────────────────────────────────────────────────────────────────────────┘

1. TOPIC
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un TOPIC = catégorie ou flux de messages

ANALOGIE : Chaîne YouTube
-> Topic "payments" = Chaîne "Payments"
-> Chaque message = une vidéo dans la chaîne
-> Abonnés = consumers

CARACTÉRISTIQUES :
-> Nom unique dans le cluster
-> Divisé en PARTITIONS (pour parallélisme)
-> Retention configurable (temps ou taille)
-> Multi-producers et multi-consumers


EXEMPLES DE TOPICS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

E-commerce :
-> orders
-> payments
-> shipments
-> user-registrations
-> product-views

IoT :
-> sensor-temperature
-> sensor-humidity
-> device-status

Logs :
-> app-logs
-> access-logs
-> error-logs


CONVENTION DE NOMMAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Bonnes pratiques :
[OK] Utiliser des noms descriptifs
[OK] Utiliser des tirets ou underscores
[OK] Éviter les espaces
[OK] Préfixer par domaine/environnement si nécessaire

Exemples :
-> prod.orders
-> staging.user-events
-> analytics_clickstream
-> payments-v2


2. PARTITION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Une PARTITION = sous-division ordonnée d'un topic

POURQUOI DES PARTITIONS ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. PARALLÉLISME : plusieurs consumers en parallèle
2. SCALABILITÉ : distribuer la charge sur plusieurs brokers
3. ORDRE GARANTI : dans une partition (pas entre partitions)

STRUCTURE D'UN TOPIC AVEC 3 PARTITIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic: "orders"

Partition 0:  [Msg0][Msg3][Msg6][Msg9][Msg12]...
              Offset: 0    1    2    3    4

Partition 1:  [Msg1][Msg4][Msg7][Msg10][Msg13]...
              Offset: 0    1    2    3     4

Partition 2:  [Msg2][Msg5][Msg8][Msg11][Msg14]...
              Offset: 0    1    2    3     4

-> Chaque partition = log INDÉPENDANT
-> Chaque partition a ses propres offsets (commence à 0)
-> Messages dans UNE partition sont ORDONNÉS
-> Ordre GLOBAL non garanti entre partitions


DISTRIBUTION DES MESSAGES DANS LES PARTITIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Quand un producer envoie un message, Kafka décide de la partition :

CAS 1 : Message AVEC clé (key)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Kafka utilise : hash(key) % nombre_de_partitions
-> MÊME clé -> TOUJOURS la MÊME partition
-> Garantit l'ORDRE pour tous les messages d'une même clé

Exemple :
Key = "user-123"
hash("user-123") = 456789
456789 % 3 = 0  ->  Partition 0

Tous les messages avec key="user-123" iront en Partition 0

CAS 2 : Message SANS clé (key = null)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Distribution ROUND-ROBIN ou STICKY (depuis Kafka 2.4)
-> Pas d'ordre garanti globalement

CAS 3 : Partition EXPLICITE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Producer spécifie directement la partition
-> Contrôle total mais plus complexe


EXEMPLE CONCRET : E-COMMERCE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic: "orders" (3 partitions)
Key: userId

Message 1: { userId: "u-001", orderId: "o-100", amount: 50 }
-> hash("u-001") % 3 = 1 -> Partition 1

Message 2: { userId: "u-002", orderId: "o-101", amount: 75 }
-> hash("u-002") % 3 = 2 -> Partition 2

Message 3: { userId: "u-001", orderId: "o-102", amount: 30 }
-> hash("u-001") % 3 = 1 -> Partition 1 (même utilisateur, même partition [OK])

AVANTAGE :
Toutes les commandes d'un utilisateur sont dans la MÊME partition
-> Ordre garanti des commandes par utilisateur
-> Consumer peut traiter les commandes d'un utilisateur en ordre


COMBIEN DE PARTITIONS ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÈGLE GÉNÉRALE :
-> Nombre de partitions = facteur de parallélisme souhaité

CALCUL :
Nombre de partitions ≈ max(T / P, T / C)

Où :
T = Débit cible (throughput)
P = Débit par partition en production
C = Débit par partition en consommation

EXEMPLE :
Débit cible : 100 MB/s
Débit producer par partition : 10 MB/s
Débit consumer par partition : 10 MB/s

Partitions = max(100/10, 100/10) = 10 partitions

CONSIDÉRATIONS :
[OK] Plus de partitions = plus de parallélisme
[X] Plus de partitions = plus de overhead (métadonnées)
[X] Plus de partitions = élection de leader plus longue en cas de panne

RECOMMENDATIONS :
-> Petit volume : 3-6 partitions
-> Moyen volume : 10-30 partitions
-> Gros volume : 50-100 partitions
-> Très gros volume : 100-1000 partitions

[ATTENTION]  IMPORTANT : Impossible de RÉDUIRE le nombre de partitions
    (on peut seulement augmenter)


3. OFFSET
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un OFFSET = identifiant de position d'un message dans une partition

Partition 0:
┌─────────────────────────────────────────────────────────────────┐
│ Offset:  0      1      2      3      4      5      6      7     │
│ Msg:   [Msg0] [Msg1] [Msg2] [Msg3] [Msg4] [Msg5] [Msg6] [Msg7]  │
└─────────────────────────────────────────────────────────────────┘

CARACTÉRISTIQUES :
-> Entier incrémental (0, 1, 2, 3, ...)
-> Unique dans UNE partition (pas globalement)
-> Immuable (une fois assigné, ne change jamais)
-> Géré par chaque consumer (pas par Kafka)

CONSUMER OFFSET :
Chaque consumer group garde la trace de son dernier offset lu :

Consumer Group "analytics":
-> Partition 0, Offset: 1250
-> Partition 1, Offset: 980
-> Partition 2, Offset: 1100

Consumer Group "billing":
-> Partition 0, Offset: 3500
-> Partition 1, Offset: 3200
-> Partition 2, Offset: 3450

-> Chaque group lit INDÉPENDAMMENT
-> "analytics" peut être à l'offset 1250 pendant que "billing" est à 3500


COMMIT D'OFFSET :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Quand un consumer lit un message, il doit "commit" son offset :

1. AUTO COMMIT (par défaut) :
-> Kafka commit automatiquement toutes les 5 secondes
-> Risque de duplication si crash avant commit
-> Simple mais moins de contrôle

2. MANUAL COMMIT (recommandé en prod) :
-> Consumer commit explicitement après traitement
-> Plus de contrôle
-> Garanties plus fortes


4. BROKER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un BROKER = serveur Kafka qui stocke les messages

┌─────────────────────────────────────────────────────────────────┐
│                        KAFKA BROKER                             │
│                                                                 │
│  ┌────────────────────────────────────────────────────────────┐ │
│  │  TOPICS & PARTITIONS                                       │ │
│  │                                                            │ │
│  │  Topic: orders                                             │ │
│  │  ├─ Partition 0: [Msg0][Msg1][Msg2]...                     │ │
│  │  ├─ Partition 1: [Msg0][Msg1][Msg2]...                     │ │
│  │  └─ Partition 2: [Msg0][Msg1][Msg2]...                     │ │
│  │                                                            │ │
│  │  Topic: payments                                           │ │
│  │  ├─ Partition 0: [Msg0][Msg1][Msg2]...                     │ │
│  │  └─ Partition 1: [Msg0][Msg1][Msg2]...                     │ │
│  └────────────────────────────────────────────────────────────┘ │
│                                                                 │
│  ┌────────────────────────────────────────────────────────────┐ │
│  │  STOCKAGE SUR DISQUE                                       │ │
│  │  /var/kafka-logs/                                          │ │
│  │  ├─ orders-0/                                              │ │
│  │  │  ├─ 00000000000000000000.log                            │ │
│  │  │  ├─ 00000000000000000000.index                          │ │
│  │  │  └─ 00000000000000000000.timeindex                      │ │
│  │  └─ orders-1/                                              │ │
│  └────────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘

RÔLES D'UN BROKER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Recevoir les messages des producers
2. Stocker les messages sur disque
3. Servir les messages aux consumers
4. Répliquer les partitions sur d'autres brokers
5. Gérer les élections de leader


5. CLUSTER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un CLUSTER = ensemble de brokers Kafka

┌──────────────────────────────────────────────────────────────┐
│                        KAFKA CLUSTER                         │
│                                                              │
│  ┌───────────────┐  ┌───────────────┐  ┌───────────────┐     │
│  │   BROKER 1    │  │   BROKER 2    │  │   BROKER 3    │     │
│  │   (ID: 101)   │  │   (ID: 102)   │  │   (ID: 103)   │     │
│  │               │  │               │  │               │     │
│  │ orders-0 (L)  │  │ orders-0 (F)  │  │ orders-1 (L)  │     │
│  │ orders-1 (F)  │  │ orders-1 (F)  │  │ orders-2 (L)  │     │
│  │ orders-2 (F)  │  │ payments-0 (L)│  │ payments-1(L) │     │
│  └───────────────┘  └───────────────┘  └───────────────┘     │
│          │                  │                  │             │
│          └──────────────────┴──────────────────┘             │
│                             │                                │
│                  ┌──────────────────┐                        │
│                  │   ZOOKEEPER      │                        │
│                  │   ou KRaft       │                        │
│                  └──────────────────┘                        │
└──────────────────────────────────────────────────────────────┘

L = Leader
F = Follower (replica)

AVANTAGES DU CLUSTER :
[OK] Haute disponibilité (si un broker tombe, autres continuent)
[OK] Scalabilité (ajouter des brokers pour plus de débit)
[OK] Réplication (copies des données sur plusieurs brokers)


6. LEADER ET FOLLOWERS (RÉPLICATION)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque PARTITION a :
-> 1 LEADER (reçoit toutes les lectures/écritures)
-> N FOLLOWERS (répliques, synchronisées avec le leader)

RÉPLICATION FACTOR = 3 :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic: "orders", Partition 0

┌──────────────┐       ┌──────────────┐       ┌──────────────┐
│  BROKER 1    │       │  BROKER 2    │       │  BROKER 3    │
│              │       │              │       │              │
│ orders-0     │       │ orders-0     │       │ orders-0     │
│ (LEADER) [OK]  │ ────-> │ (FOLLOWER)   │ ────-> │ (FOLLOWER)   │
│              │       │              │       │              │
│ [Msg0]       │       │ [Msg0]       │       │ [Msg0]       │
│ [Msg1]       │       │ [Msg1]       │       │ [Msg1]       │
│ [Msg2]       │       │ [Msg2]       │       │ [Msg2]       │
└──────────────┘       └──────────────┘       └──────────────┘
      [BLACK_UP-POINTING_TRIANGLE]
      │
   Producer
   Consumer

FONCTIONNEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Producer envoie message au LEADER uniquement
2. Leader écrit le message dans son log
3. Followers PULL (tirent) les messages depuis le leader
4. Followers écrivent dans leur log
5. Followers envoient ACK au leader
6. Consumer lit depuis le LEADER uniquement (par défaut)

SI LE LEADER TOMBE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. ZooKeeper (ou KRaft) détecte la panne
2. Élection automatique d'un nouveau leader parmi les followers
3. Le nouveau leader reprend les lectures/écritures
4. Downtime : quelques secondes seulement

ISR (In-Sync Replicas) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ISR = Liste des répliques qui sont SYNCHRONISÉES avec le leader

Exemple :
Partition orders-0
-> Leader: Broker 1
-> ISR: [Broker 1, Broker 2, Broker 3]

Si Broker 3 prend du retard :
-> ISR: [Broker 1, Broker 2]
-> Broker 3 hors ISR (pas éligible comme leader)

ACKS (Acknowledgments) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Producer peut configurer le niveau de garantie :

acks = 0 (fire and forget)
-> Producer n'attend PAS d'ACK
-> Ultra rapide
-> Risque de perte de données [X]

acks = 1 (leader acknowledgment)
-> Leader ACK immédiatement après écriture locale
-> Rapide
-> Risque de perte si leader crash avant réplication [X]

acks = all (ou -1) (all ISR acknowledgment)
-> Leader attend ACK de TOUS les followers ISR
-> Plus lent
-> Garantie maximale [OK]
-> Recommandé pour données critiques


7. ZOOKEEPER vs KRaft
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ZOOKEEPER (traditionnel, < Kafka 3.3) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Rôle de ZooKeeper :
-> Coordination du cluster
-> Élection des leaders
-> Stockage des métadonnées (topics, partitions, ISR)
-> Gestion de la configuration

Architecture :
┌─────────────┐
│  Kafka      │
│  Cluster    │
└──────┬──────┘
       │ dépend
       [BLACK_DOWN-POINTING_TRIANGLE]
┌─────────────┐
│  ZooKeeper  │
│  Ensemble   │
└─────────────┘

Problèmes :
[X] Complexité (2 systèmes à gérer)
[X] Latence supplémentaire
[X] Point de défaillance unique


KRaft (Kafka Raft) - Nouveau mode (depuis Kafka 2.8, stable 3.3+) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Rôle de KRaft :
-> Remplace ZooKeeper complètement
-> Consensus Raft intégré dans Kafka
-> Métadonnées gérées par des brokers "controllers"

Architecture :
┌────────────────────────────────────┐
│  Kafka Cluster (KRaft mode)        │
│                                    │
│  ┌─────────┐  ┌─────────┐          │
│  │ Broker  │  │ Broker  │          │
│  │(regular)│  │(control)│          │
│  └─────────┘  └─────────┘          │
└────────────────────────────────────┘

Avantages :
[OK] Plus simple (un seul système)
[OK] Plus rapide (pas de hop ZooKeeper)
[OK] Meilleure scalabilité (millions de partitions)
[OK] Déploiement simplifié

STATUT :
-> Production-ready depuis Kafka 3.3 (2022)
-> ZooKeeper sera déprécié dans Kafka 4.0


RÉSUMÉ ARCHITECTURE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

HIÉRARCHIE :
Cluster
  └─ Brokers (1 à N)
      └─ Topics (1 à M)
          └─ Partitions (1 à P)
              └─ Messages (offsets 0, 1, 2, ...)
                  └─ Leader + Followers (réplication)

FLUX DE DONNÉES :
1. Producer -> envoie message -> Leader de la partition
2. Leader -> stocke -> réplique aux Followers
3. Consumer -> lit depuis -> Leader
4. Consumer -> commit offset -> Kafka (__consumer_offsets topic)

KEY TAKEAWAYS :
[OK] Topic = catégorie de messages
[OK] Partition = parallélisme et scalabilité
[OK] Offset = position d'un message
[OK] Broker = serveur de stockage
[OK] Cluster = groupe de brokers
[OK] Leader/Followers = haute disponibilité
[OK] ISR = répliques synchronisées
[OK] acks = garantie de durabilité
"""


# [OK] PARTIE 3 : INSTALLATION ET CONFIGURATION

"""
┌────────────────────────────────────────────────────────────────────────┐
│              INSTALLATION KAFKA (TOUS LES OS)                          │
└────────────────────────────────────────────────────────────────────────┘

PRÉREQUIS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Java 11+ (JDK) installé
-> 4 GB RAM minimum (8+ GB recommandé)
-> Plusieurs GB d'espace disque


MÉTHODES D'INSTALLATION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Binaires officiels (Apache Kafka)
2. Docker (recommandé pour développement)
3. Confluent Platform (distribution commerciale)
4. Cloud managé (AWS MSK, Confluent Cloud, etc.)


MÉTHODE 1 : BINAIRES OFFICIELS (LINUX/MAC)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# 1. Télécharger Kafka
cd /opt
wget https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz
tar -xzf kafka_2.13-3.6.1.tgz
cd kafka_2.13-3.6.1

# 2. Démarrer Kafka avec KRaft (sans ZooKeeper)
# Générer un UUID pour le cluster
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"

# Formater le répertoire de logs
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties

# Démarrer Kafka
bin/kafka-server-start.sh config/kraft/server.properties

# Kafka démarre sur port 9092 par défaut


# 3. (ALTERNATIF) Démarrer avec ZooKeeper (ancien mode)
# Terminal 1 : Démarrer ZooKeeper
bin/zookeeper-server-start.sh config/zookeeper.properties

# Terminal 2 : Démarrer Kafka
bin/kafka-server-start.sh config/server.properties


MÉTHODE 2 : DOCKER (RECOMMANDÉ POUR DÉVELOPPEMENT)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# docker-compose.yml (KRaft mode)
version: '3.8'

services:
  kafka:
    image: apache/kafka:3.6.1
    container_name: kafka
    ports:
      - "9092:9092"
    environment:
      # KRaft mode configuration
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
      CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
    volumes:
      - kafka-data:/tmp/kraft-combined-logs

volumes:
  kafka-data:

# Démarrer
docker-compose up -d

# Vérifier les logs
docker-compose logs -f kafka


# docker-compose.yml (avec ZooKeeper, plus complet)
version: '3.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - "2181:2181"

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "9093:9093"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_HOST://localhost:9093
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
    volumes:
      - kafka-data:/var/lib/kafka/data

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    container_name: kafka-ui
    depends_on:
      - kafka
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181

volumes:
  kafka-data:


CLUSTER KAFKA (3 brokers) AVEC DOCKER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# docker-compose-cluster.yml
version: '3.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka1:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka1:19092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2

  kafka2:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9093:9093"
    environment:
      KAFKA_BROKER_ID: 2
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka2:19093,PLAINTEXT_HOST://localhost:9093
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT

  kafka3:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9094:9094"
    environment:
      KAFKA_BROKER_ID: 3
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka3:19094,PLAINTEXT_HOST://localhost:9094
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT


CONFIGURATION KAFKA :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Fichier : config/server.properties (ou config/kraft/server.properties)

# ═══════════════════════════════════════════════════════════════════
# BROKER
# ═══════════════════════════════════════════════════════════════════
broker.id=1
# Identifiant unique du broker dans le cluster

# ═══════════════════════════════════════════════════════════════════
# LISTENERS
# ═══════════════════════════════════════════════════════════════════
listeners=PLAINTEXT://0.0.0.0:9092
# Sur quelles interfaces écouter

advertised.listeners=PLAINTEXT://localhost:9092
# Adresse communiquée aux clients (producers/consumers)

listener.security.protocol.map=PLAINTEXT:PLAINTEXT
# Mapping protocol -> sécurité

# ═══════════════════════════════════════════════════════════════════
# STOCKAGE
# ═══════════════════════════════════════════════════════════════════
log.dirs=/var/kafka-logs
# Répertoire de stockage des logs (messages)

num.network.threads=3
# Threads réseau pour I/O

num.io.threads=8
# Threads pour écriture/lecture disque

socket.send.buffer.bytes=102400
# Buffer TCP d'envoi (100 KB)

socket.receive.buffer.bytes=102400
# Buffer TCP de réception (100 KB)

socket.request.max.bytes=104857600
# Taille max d'une requête (100 MB)

# ═══════════════════════════════════════════════════════════════════
# TOPICS
# ═══════════════════════════════════════════════════════════════════
num.partitions=3
# Nombre de partitions par défaut pour les nouveaux topics

default.replication.factor=3
# Facteur de réplication par défaut

min.insync.replicas=2
# Nombre minimum de répliques ISR pour acks=all

auto.create.topics.enable=true
# Création automatique des topics (false en production recommandé)

# ═══════════════════════════════════════════════════════════════════
# RÉTENTION
# ═══════════════════════════════════════════════════════════════════
log.retention.hours=168
# Rétention par temps (168h = 7 jours)

log.retention.bytes=-1
# Rétention par taille (-1 = infini)

log.segment.bytes=1073741824
# Taille max d'un segment (1 GB)

log.retention.check.interval.ms=300000
# Fréquence de vérification de la rétention (5 min)

log.segment.delete.delay.ms=60000
# Délai avant suppression d'un segment (60 sec)

# ═══════════════════════════════════════════════════════════════════
# ZOOKEEPER (si mode ZooKeeper)
# ═══════════════════════════════════════════════════════════════════
zookeeper.connect=localhost:2181
# Adresse du ZooKeeper ensemble

zookeeper.connection.timeout.ms=18000
# Timeout de connexion à ZooKeeper

# ═══════════════════════════════════════════════════════════════════
# KRAFT (si mode KRaft)
# ═══════════════════════════════════════════════════════════════════
process.roles=broker,controller
# Rôles du nœud (broker, controller, ou les deux)

node.id=1
# ID unique du nœud

controller.quorum.voters=1@localhost:9093
# Liste des voters du quorum (format: id@host:port)

controller.listener.names=CONTROLLER
# Nom du listener pour communication controller

# ═══════════════════════════════════════════════════════════════════
# PERFORMANCE
# ═══════════════════════════════════════════════════════════════════
num.replica.fetchers=4
# Threads de réplication par broker

replica.fetch.max.bytes=1048576
# Taille max d'un fetch de réplication (1 MB)

replica.lag.time.max.ms=30000
# Lag max avant qu'un follower soit retiré de l'ISR

compression.type=producer
# Type de compression (none, gzip, snappy, lz4, zstd, producer)

message.max.bytes=1048576
# Taille max d'un message (1 MB)

# ═══════════════════════════════════════════════════════════════════
# GROUPE DE CONSOMMATEURS
# ═══════════════════════════════════════════════════════════════════
group.initial.rebalance.delay.ms=3000
# Délai avant premier rebalancing (3 sec)

offsets.topic.replication.factor=3
# Réplication du topic __consumer_offsets

transaction.state.log.replication.factor=3
# Réplication du topic des transactions

transaction.state.log.min.isr=2
# ISR minimum pour transactions

# ═══════════════════════════════════════════════════════════════════
# JVM (dans kafka-server-start.sh ou variables d'environnement)
# ═══════════════════════════════════════════════════════════════════
export KAFKA_HEAP_OPTS="-Xmx6G -Xms6G"
# Heap JVM (6 GB)

export KAFKA_JVM_PERFORMANCE_OPTS="-XX:+UseG1GC -XX:MaxGCPauseMillis=20"
# GC optimisé


COMMANDES CLI ESSENTIELLES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# 1. CRÉER UN TOPIC
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create \
  --topic orders \
  --partitions 3 \
  --replication-factor 3 \
  --config retention.ms=604800000

# 2. LISTER LES TOPICS
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list

# 3. DÉCRIRE UN TOPIC
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe \
  --topic orders

Output:
Topic: orders   PartitionCount: 3   ReplicationFactor: 3
  Topic: orders Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
  Topic: orders Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
  Topic: orders Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2

# 4. MODIFIER UN TOPIC (ajouter des partitions)
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --alter \
  --topic orders \
  --partitions 5

[ATTENTION] On peut AUGMENTER le nombre de partitions, mais PAS le réduire !

# 5. SUPPRIMER UN TOPIC
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --delete \
  --topic orders

# 6. PRODUIRE DES MESSAGES (console)
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
  --topic orders

> {"orderId": "o-001", "amount": 100}
> {"orderId": "o-002", "amount": 150}

# Produire avec une clé
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
  --topic orders \
  --property "parse.key=true" \
  --property "key.separator=:"

> user-123:{"orderId": "o-001", "amount": 100}
> user-456:{"orderId": "o-002", "amount": 150}

# 7. CONSOMMER DES MESSAGES (console)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders \
  --from-beginning

# Consommer avec clés et timestamps
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders \
  --from-beginning \
  --property print.key=true \
  --property print.timestamp=true

# 8. LISTER LES CONSUMER GROUPS
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list

# 9. DÉCRIRE UN CONSUMER GROUP
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe \
  --group my-consumer-group

Output:
GROUP           TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
my-group        orders    0          150             150             0
my-group        orders    1          120             125             5
my-group        orders    2          180             180             0

LAG = nombre de messages non encore consommés

# 10. RESET OFFSET D'UN CONSUMER GROUP
# Reset au début
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --topic orders \
  --reset-offsets --to-earliest \
  --execute

# Reset à une date spécifique
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --topic orders \
  --reset-offsets --to-datetime 2024-12-01T00:00:00.000 \
  --execute

# 11. VÉRIFIER LES PERFORMANCES
bin/kafka-producer-perf-test.sh --topic test-perf \
  --num-records 100000 \
  --record-size 1000 \
  --throughput -1 \
  --producer-props bootstrap.servers=localhost:9092

bin/kafka-consumer-perf-test.sh --bootstrap-server localhost:9092 \
  --topic test-perf \
  --messages 100000 \
  --threads 1


TOOLS DE GESTION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. KAFKA UI (Web Interface)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> https://github.com/provectus/kafka-ui
-> Interface web pour gérer Kafka
-> Voir topics, partitions, messages
-> Gérer consumer groups
-> Monitoring

docker run -p 8080:8080 \
  -e KAFKA_CLUSTERS_0_NAME=local \
  -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=host.docker.internal:9092 \
  provectuslabs/kafka-ui:latest

-> Accès : http://localhost:8080

2. CONFLUENT CONTROL CENTER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Partie de Confluent Platform
-> Interface professionnelle
-> Monitoring avancé
-> Alertes

3. KAFDROP
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> https://github.com/obsidiandynamics/kafdrop
-> Alternative légère à Kafka UI

docker run -p 9000:9000 \
  -e KAFKA_BROKERCONNECT=host.docker.internal:9092 \
  obsidiandynamics/kafdrop:latest
"""


# [OK] PARTIE 4 : PRODUCERS - ÉCRIRE DES MESSAGES

"""
┌────────────────────────────────────────────────────────────────────────┐
│              KAFKA PRODUCERS : PUBLIER DES MESSAGES                    │
└────────────────────────────────────────────────────────────────────────┘

QU'EST-CE QU'UN PRODUCER ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un PRODUCER = application qui PUBLIE des messages dans un topic Kafka

┌──────────────┐                    ┌────────────────┐
│  PRODUCER    │ ─── messages ────-> │  KAFKA BROKER  │
│  (App)       │                    │  (Topic)       │
└──────────────┘                    └────────────────┘


EXEMPLE JAVA (Kafka Client)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// pom.xml
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.6.1</version>
</dependency>

// SimpleProducer.java
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class SimpleProducer {
    public static void main(String[] args) {
        
        // 1. CONFIGURATION DU PRODUCER
        Properties props = new Properties();
        
        // Adresse du broker Kafka
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        
        // Sérialisation de la clé et de la valeur
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
                  StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
                  StringSerializer.class.getName());
        
        // Acknowledgment (garantie de livraison)
        props.put(ProducerConfig.ACKS_CONFIG, "all");  // Attendre tous les ISR
        
        // Retries en cas d'erreur
        props.put(ProducerConfig.RETRIES_CONFIG, 3);
        
        // 2. CRÉER LE PRODUCER
        KafkaProducer<String, String> producer = new KafkaProducer<>(props);
        
        try {
            // 3. ENVOYER DES MESSAGES
            for (int i = 0; i < 10; i++) {
                String key = "user-" + i;
                String value = "{\"orderId\": \"o-" + i + "\", \"amount\": " + (100 + i * 10) + "}";
                
                ProducerRecord<String, String> record = 
                    new ProducerRecord<>("orders", key, value);
                
                // Fire and forget (asynchrone)
                producer.send(record);
                
                System.out.println("Message envoyé: " + value);
            }
            
        } finally {
            // 4. FERMER LE PRODUCER
            producer.close();
        }
    }
}


PRODUCER AVEC CALLBACK (Synchronisation)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

for (int i = 0; i < 10; i++) {
    ProducerRecord<String, String> record = 
        new ProducerRecord<>("orders", "user-" + i, "...");
    
    // Envoi avec callback
    producer.send(record, new Callback() {
        @Override
        public void onCompletion(RecordMetadata metadata, Exception exception) {
            if (exception == null) {
                // Succès
                System.out.println("Message envoyé avec succès !");
                System.out.println("Topic: " + metadata.topic());
                System.out.println("Partition: " + metadata.partition());
                System.out.println("Offset: " + metadata.offset());
                System.out.println("Timestamp: " + metadata.timestamp());
            } else {
                // Erreur
                System.err.println("Erreur lors de l'envoi: " + exception.getMessage());
            }
        }
    });
}

// Attendre que tous les messages soient envoyés
producer.flush();


PRODUCER SYNCHRONE (Attendre confirmation)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

try {
    RecordMetadata metadata = producer.send(record).get();  // .get() = bloquant
    System.out.println("Offset: " + metadata.offset());
} catch (InterruptedException | ExecutionException e) {
    System.err.println("Erreur: " + e.getMessage());
}


EXEMPLE PYTHON (kafka-python)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Installation
pip install kafka-python

# simple_producer.py
from kafka import KafkaProducer
import json

# 1. Créer le producer
producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8'),
    key_serializer=lambda k: k.encode('utf-8') if k else None,
    acks='all',
    retries=3
)

# 2. Envoyer des messages
for i in range(10):
    key = f"user-{i}"
    value = {
        "orderId": f"o-{i}",
        "amount": 100 + i * 10
    }
    
    # Fire and forget
    producer.send('orders', key=key, value=value)
    print(f"Message envoyé: {value}")

# 3. Attendre l'envoi de tous les messages
producer.flush()

# 4. Fermer le producer
producer.close()


# Producer avec callback
def on_send_success(record_metadata):
    print(f"Topic: {record_metadata.topic}")
    print(f"Partition: {record_metadata.partition}")
    print(f"Offset: {record_metadata.offset}")

def on_send_error(excp):
    print(f"Erreur: {excp}")

# Envoi avec callback
future = producer.send('orders', key=key, value=value)
future.add_callback(on_send_success)
future.add_errback(on_send_error)

# Envoi synchrone
try:
    record_metadata = future.get(timeout=10)
    print(f"Offset: {record_metadata.offset}")
except Exception as e:
    print(f"Erreur: {e}")


EXEMPLE NODE.JS (kafkajs)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// Installation
npm install kafkajs

// simple-producer.js
const { Kafka } = require('kafkajs');

// 1. Créer l'instance Kafka
const kafka = new Kafka({
  clientId: 'my-app',
  brokers: ['localhost:9092']
});

// 2. Créer le producer
const producer = kafka.producer();

async function run() {
  // 3. Se connecter
  await producer.connect();
  console.log('Producer connecté [OK]');
  
  // 4. Envoyer des messages
  for (let i = 0; i < 10; i++) {
    const message = {
      key: `user-${i}`,
      value: JSON.stringify({
        orderId: `o-${i}`,
        amount: 100 + i * 10
      })
    };
    
    await producer.send({
      topic: 'orders',
      messages: [message]
    });
    
    console.log(`Message envoyé: ${message.value}`);
  }
  
  // 5. Déconnecter
  await producer.disconnect();
  console.log('Producer déconnecté [OK]');
}

run().catch(console.error);


// Envoi en batch (plus performant)
const messages = [];
for (let i = 0; i < 100; i++) {
  messages.push({
    key: `user-${i}`,
    value: JSON.stringify({ orderId: `o-${i}`, amount: 100 })
  });
}

await producer.send({
  topic: 'orders',
  messages: messages
});


CONFIGURATION PRODUCER IMPORTANTE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. bootstrap.servers
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Liste des brokers Kafka

props.put("bootstrap.servers", "localhost:9092,localhost:9093");

-> Producer contacte ces brokers pour obtenir les métadonnées du cluster
-> Pas besoin de lister TOUS les brokers (2-3 suffisent)


2. key.serializer & value.serializer
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Comment convertir clé/valeur en bytes

Serializers disponibles :
-> StringSerializer
-> IntegerSerializer
-> LongSerializer
-> ByteArraySerializer
-> Avro/Protobuf (avec Schema Registry)

props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");


3. acks (Acknowledgment)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Niveau de garantie de livraison

acks=0 (No acknowledgment)
-> Producer n'attend PAS d'ACK du broker
-> Très rapide
-> Risque de perte de données [X]
-> Use case : Metrics non critiques, logs

acks=1 (Leader acknowledgment)
-> Leader ACK après écriture locale
-> Rapide
-> Risque de perte si leader crash avant réplication [X]
-> Use case : Logs applicatifs

acks=all (ou -1) (All ISR acknowledgment)
-> Leader attend ACK de TOUS les followers ISR
-> Plus lent
-> Garantie maximale [OK]
-> Use case : Transactions financières, données critiques

props.put("acks", "all");


4. retries & retry.backoff.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Nombre de tentatives en cas d'erreur

props.put("retries", 3);
// Réessayer 3 fois en cas d'erreur transiente

props.put("retry.backoff.ms", 100);
// Attendre 100ms entre chaque retry


5. compression.type
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Type de compression des messages

Options :
-> none (pas de compression)
-> gzip (bonne compression, CPU moyen)
-> snappy (compression moyenne, CPU faible)
-> lz4 (compression moyenne, très rapide)
-> zstd (meilleure compression, CPU moyen) - recommandé

props.put("compression.type", "snappy");

AVANTAGES :
[OK] Réduction de la bande passante réseau
[OK] Réduction de l'espace disque
[OK] Meilleur débit global

INCONVÉNIENTS :
[X] CPU supplémentaire
[X] Latence légèrement supérieure


6. batch.size & linger.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Batching pour performance

batch.size (bytes)
-> Taille maximale d'un batch
-> Par défaut : 16384 (16 KB)

props.put("batch.size", 32768);  // 32 KB

linger.ms (millisecondes)
-> Temps d'attente avant d'envoyer un batch incomplet
-> Par défaut : 0 (envoyer immédiatement)

props.put("linger.ms", 10);  // Attendre 10ms pour remplir le batch

TRADE-OFF :
-> linger.ms = 0 -> Latence minimale, débit faible
-> linger.ms > 0 -> Latence augmente, débit élevé


7. buffer.memory
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Mémoire totale pour le buffering

props.put("buffer.memory", 33554432);  // 32 MB

Si buffer plein :
-> Producer bloque (max.block.ms)
-> Ou lève une exception


8. max.in.flight.requests.per.connection
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Nombre de requêtes en vol sans attendre ACK

props.put("max.in.flight.requests.per.connection", 5);

IMPORTANT POUR L'ORDRE :
-> Si > 1 ET retries > 0 -> Risque de réordonnancement
-> Pour garantir l'ordre strict : mettre à 1 (mais débit réduit)
-> Depuis Kafka 1.0 : enable.idempotence=true résout ce problème


9. enable.idempotence
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Garantir l'idempotence (pas de duplicata)

props.put("enable.idempotence", true);

EFFETS :
[OK] Pas de message dupliqué (exactly-once semantics)
[OK] Ordre garanti même avec retries
[OK] acks=all automatiquement
[OK] max.in.flight.requests <= 5

RECOMMANDÉ EN PRODUCTION !


PARTITIONNEUR (Partitioner)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Le Partitioner décide vers quelle partition envoyer un message

STRATÉGIES PAR DÉFAUT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Message AVEC clé :
   partition = hash(key) % num_partitions

2. Message SANS clé :
   Round-robin ou Sticky partitioning (depuis Kafka 2.4)


CUSTOM PARTITIONER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// CustomPartitioner.java
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;

import java.util.Map;

public class CustomPartitioner implements Partitioner {
    
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        
        int numPartitions = cluster.partitionCountForTopic(topic);
        
        // Exemple : Utilisateurs VIP vont en partition 0
        if (key != null && key.toString().startsWith("vip-")) {
            return 0;
        }
        
        // Autres : distribution normale
        return Math.abs(key.hashCode()) % numPartitions;
    }
    
    @Override
    public void close() {
        // Nettoyage si nécessaire
    }
    
    @Override
    public void configure(Map<String, ?> configs) {
        // Configuration si nécessaire
    }
}

// Utilisation
props.put("partitioner.class", CustomPartitioner.class.getName());


GESTION D'ERREURS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TYPES D'ERREURS :

1. ERREURS RETRIABLES (transientes)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> NotLeaderForPartitionException
-> NetworkException
-> TimeoutException

-> Producer RETRY automatiquement

2. ERREURS NON-RETRIABLES (permanentes)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> RecordTooLargeException (message > max.message.bytes)
-> SerializationException (erreur de sérialisation)
-> InvalidTopicException

-> Producer lève une exception immédiatement

producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        if (exception instanceof RetriableException) {
            // Erreur retriable, Kafka réessaiera
            log.warn("Erreur retriable: {}", exception.getMessage());
        } else {
            // Erreur non-retriable, action nécessaire
            log.error("Erreur permanente: {}", exception.getMessage());
            // Envoyer vers Dead Letter Queue
            sendToDeadLetterQueue(record);
        }
    }
});


BONNES PRATIQUES PRODUCER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] TOUJOURS utiliser enable.idempotence=true
[OK] Utiliser acks=all pour données critiques
[OK] Fermer le producer proprement (close())
[OK] Utiliser compression (snappy ou zstd)
[OK] Batching avec linger.ms > 0 pour débit élevé
[OK] Utiliser des clés pour garantir l'ordre
[OK] Gérer les erreurs avec callbacks
[OK] Monitorer les métriques (record-send-rate, compression-rate)
[OK] Tester les retries et timeouts
[OK] Un seul producer par application (thread-safe)

[X] Ne PAS créer un producer par message (overhead)
[X] Ne PAS oublier flush() avant close()
[X] Ne PAS utiliser acks=0 pour données importantes
[X] Ne PAS ignorer les exceptions
[X] Ne PAS bloquer le thread du callback
"""


Je continue le guide Kafka avec les parties suivantes :

```txt
# [OK] PARTIE 5 : CONSUMERS - LIRE DES MESSAGES

"""
┌────────────────────────────────────────────────────────────────────────┐
│              KAFKA CONSUMERS : CONSOMMER DES MESSAGES                  │
└────────────────────────────────────────────────────────────────────────┘

QU'EST-CE QU'UN CONSUMER ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un CONSUMER = application qui LIT des messages depuis un topic Kafka

┌────────────────┐                    ┌──────────────┐
│  KAFKA BROKER  │ ─── messages ────-> │  CONSUMER    │
│  (Topic)       │                    │  (App)       │
└────────────────┘                    └──────────────┘


EXEMPLE JAVA (Kafka Client)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// SimpleConsumer.java
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class SimpleConsumer {
    public static void main(String[] args) {
        
        // 1. CONFIGURATION DU CONSUMER
        Properties props = new Properties();
        
        // Adresse du broker Kafka
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        
        // Group ID (identifiant du consumer group)
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-consumer-group");
        
        // Désérialisation de la clé et de la valeur
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, 
                  StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, 
                  StringDeserializer.class.getName());
        
        // Auto commit des offsets
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
        props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "5000");
        
        // Offset initial si pas de commit précédent
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        // Options: earliest, latest, none
        
        // 2. CRÉER LE CONSUMER
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
        
        // 3. SOUSCRIRE À UN OU PLUSIEURS TOPICS
        consumer.subscribe(Collections.singletonList("orders"));
        
        try {
            // 4. BOUCLE DE CONSOMMATION
            while (true) {
                // Poll : récupérer des messages (timeout 100ms)
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                
                // 5. TRAITER LES MESSAGES
                for (ConsumerRecord<String, String> record : records) {
                    System.out.println("Message reçu:");
                    System.out.println("  Topic: " + record.topic());
                    System.out.println("  Partition: " + record.partition());
                    System.out.println("  Offset: " + record.offset());
                    System.out.println("  Key: " + record.key());
                    System.out.println("  Value: " + record.value());
                    System.out.println("  Timestamp: " + record.timestamp());
                    System.out.println("---");
                }
            }
        } finally {
            // 6. FERMER LE CONSUMER
            consumer.close();
        }
    }
}


CONSUMER AVEC COMMIT MANUEL (Recommandé)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "my-consumer-group");
props.put("enable.auto.commit", "false");  // [ATTENTION] Désactiver auto-commit
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("orders"));

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        try {
            // Traiter le message
            processMessage(record);
            
            // Commit synchrone après chaque message
            consumer.commitSync();
            
        } catch (Exception e) {
            // En cas d'erreur, on ne commit pas
            // Le message sera retraité
            System.err.println("Erreur de traitement: " + e.getMessage());
        }
    }
}

// Alternative : Commit après un batch
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        processMessage(record);
    }
    
    // Commit une fois tous les messages du batch traités
    consumer.commitSync();
}


EXEMPLE PYTHON (kafka-python)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# simple_consumer.py
from kafka import KafkaConsumer
import json

# 1. Créer le consumer
consumer = KafkaConsumer(
    'orders',  # Topic à consommer
    bootstrap_servers=['localhost:9092'],
    group_id='my-consumer-group',
    value_deserializer=lambda m: json.loads(m.decode('utf-8')),
    key_deserializer=lambda k: k.decode('utf-8') if k else None,
    auto_offset_reset='earliest',  # 'earliest', 'latest', 'none'
    enable_auto_commit=True,
    auto_commit_interval_ms=5000
)

# 2. Consommer les messages
print("Consumer démarré, en attente de messages...")

for message in consumer:
    print(f"Message reçu:")
    print(f"  Topic: {message.topic}")
    print(f"  Partition: {message.partition}")
    print(f"  Offset: {message.offset}")
    print(f"  Key: {message.key}")
    print(f"  Value: {message.value}")
    print(f"  Timestamp: {message.timestamp}")
    print("---")


# Consumer avec commit manuel
consumer = KafkaConsumer(
    'orders',
    bootstrap_servers=['localhost:9092'],
    group_id='my-consumer-group',
    enable_auto_commit=False,  # Désactiver auto-commit
    value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)

for message in consumer:
    try:
        # Traiter le message
        process_message(message.value)
        
        # Commit manuel
        consumer.commit()
        
    except Exception as e:
        print(f"Erreur: {e}")
        # Ne pas committer, le message sera retraité


EXEMPLE NODE.JS (kafkajs)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// simple-consumer.js
const { Kafka } = require('kafkajs');

// 1. Créer l'instance Kafka
const kafka = new Kafka({
  clientId: 'my-app',
  brokers: ['localhost:9092']
});

// 2. Créer le consumer
const consumer = kafka.consumer({ 
  groupId: 'my-consumer-group' 
});

async function run() {
  // 3. Se connecter
  await consumer.connect();
  console.log('Consumer connecté [OK]');
  
  // 4. Souscrire au topic
  await consumer.subscribe({ 
    topic: 'orders', 
    fromBeginning: true  // équivalent de auto.offset.reset=earliest
  });
  
  // 5. Consommer les messages
  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      console.log('Message reçu:');
      console.log(`  Topic: ${topic}`);
      console.log(`  Partition: ${partition}`);
      console.log(`  Offset: ${message.offset}`);
      console.log(`  Key: ${message.key?.toString()}`);
      console.log(`  Value: ${message.value.toString()}`);
      console.log(`  Timestamp: ${message.timestamp}`);
      console.log('---');
    }
  });
}

run().catch(console.error);


// Consumer avec commit manuel
await consumer.run({
  autoCommit: false,  // Désactiver auto-commit
  eachMessage: async ({ topic, partition, message }) => {
    try {
      // Traiter le message
      await processMessage(message);
      
      // Commit manuel
      await consumer.commitOffsets([
        {
          topic,
          partition,
          offset: (parseInt(message.offset) + 1).toString()
        }
      ]);
      
    } catch (error) {
      console.error('Erreur:', error);
      // Ne pas committer, le message sera retraité
    }
  }
});


CONFIGURATION CONSUMER IMPORTANTE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. bootstrap.servers
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Liste des brokers Kafka

props.put("bootstrap.servers", "localhost:9092");


2. group.id
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Identifiant du consumer group (TRÈS IMPORTANT)

props.put("group.id", "analytics-service");

-> Tous les consumers avec le MÊME group.id forment un groupe
-> Les partitions sont RÉPARTIES entre les consumers du groupe
-> Chaque message est lu par UN SEUL consumer du groupe


3. key.deserializer & value.deserializer
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Comment convertir bytes en objets

props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());


4. auto.offset.reset
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Où commencer si pas d'offset précédent

earliest
-> Lire depuis le DÉBUT du topic
-> Utile pour nouveau consumer qui veut tout l'historique

latest (défaut)
-> Lire seulement les NOUVEAUX messages
-> Ignorer l'historique

none
-> Lever une exception si pas d'offset précédent
-> Forcer un comportement explicite

props.put("auto.offset.reset", "earliest");


5. enable.auto.commit & auto.commit.interval.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Commit automatique des offsets

enable.auto.commit=true (défaut)
-> Kafka commit automatiquement les offsets
-> Toutes les auto.commit.interval.ms millisecondes

props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "5000");  // 5 secondes

RISQUE AVEC AUTO-COMMIT :
[X] Si consumer crash APRÈS poll() mais AVANT traitement
   -> Offset commité mais message PAS traité
   -> PERTE DE MESSAGE

RECOMMANDATION :
[OK] Utiliser enable.auto.commit=false
[OK] Commit manuellement APRÈS traitement réussi


6. max.poll.records
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Nombre maximum de messages retournés par poll()

props.put("max.poll.records", 500);

-> Limiter pour éviter de surcharger le consumer
-> Trade-off : petite valeur = plus de polls, grande valeur = traitement plus long


7. max.poll.interval.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Temps maximum entre deux poll()

props.put("max.poll.interval.ms", 300000);  // 5 minutes

-> Si consumer ne poll() pas pendant ce délai
   -> Considéré comme mort
   -> Rebalancing déclenché
   -> Ses partitions réassignées

IMPORTANT :
Si traitement d'un batch prend plus de max.poll.interval.ms
-> Augmenter cette valeur
-> Ou réduire max.poll.records


8. session.timeout.ms & heartbeat.interval.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Détection de panne d'un consumer

session.timeout.ms
-> Temps max sans heartbeat avant d'être considéré mort
-> Défaut : 10000 (10 secondes)

heartbeat.interval.ms
-> Fréquence d'envoi des heartbeats
-> Défaut : 3000 (3 secondes)
-> Règle : heartbeat.interval.ms < session.timeout.ms / 3

props.put("session.timeout.ms", "10000");
props.put("heartbeat.interval.ms", "3000");


9. fetch.min.bytes & fetch.max.wait.ms
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Contrôle du fetch

fetch.min.bytes
-> Minimum de bytes à attendre avant de retourner
-> Défaut : 1 byte

props.put("fetch.min.bytes", 1024);  // Attendre 1 KB

fetch.max.wait.ms
-> Temps max d'attente si fetch.min.bytes pas atteint
-> Défaut : 500ms

props.put("fetch.max.wait.ms", 500);

TRADE-OFF :
-> fetch.min.bytes élevé = moins de requêtes, latence plus élevée
-> fetch.min.bytes faible = plus de requêtes, latence faible


10. isolation.level
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Pour les transactions

read_uncommitted (défaut)
-> Lire tous les messages, même non committés

read_committed
-> Lire seulement les messages committés (transactions)

props.put("isolation.level", "read_committed");


POLL LOOP - DÉTAILS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

QUE SE PASSE-T-IL DANS poll() ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    // ...
}

ÉTAPES INTERNES :

1. Envoi de heartbeat au coordinator
   -> "Je suis vivant"

2. Si nécessaire, participation au rebalancing
   -> Réassignation des partitions

3. Si auto-commit activé, commit des offsets
   -> Toutes les auto.commit.interval.ms

4. Fetch des messages depuis les brokers
   -> Récupération des données

5. Retour des messages au code applicatif
   -> ConsumerRecords

IMPORTANT :
-> poll() DOIT être appelé régulièrement
-> Pas de poll() pendant max.poll.interval.ms = consumer mort


GESTION DES OFFSETS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AUTO-COMMIT vs MANUAL COMMIT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AUTO-COMMIT (enable.auto.commit=true) :
┌────────────────────────────────────────────────────────────────┐
│ Timeline:                                                      │
│                                                                │
│ T0: poll() retourne messages [100, 101, 102]                  │
│ T1: Traitement de message 100                                 │
│ T2: Traitement de message 101                                 │
│ T3: Kafka AUTO-COMMIT offset 103 [OK]                           │
│ T4: [IMPACT] CRASH pendant traitement de message 102                │
│                                                                │
│ Résultat : Message 102 PERDU [X]                               │
└────────────────────────────────────────────────────────────────┘

MANUAL COMMIT (enable.auto.commit=false) :
┌────────────────────────────────────────────────────────────────┐
│ Timeline:                                                      │
│                                                                │
│ T0: poll() retourne messages [100, 101, 102]                  │
│ T1: Traitement de message 100 -> commitSync() [OK]               │
│ T2: Traitement de message 101 -> commitSync() [OK]               │
│ T3: [IMPACT] CRASH pendant traitement de message 102                │
│                                                                │
│ Résultat : Message 102 sera RETRAITÉ au redémarrage [OK]        │
└────────────────────────────────────────────────────────────────┘


COMMIT SYNCHRONE vs ASYNCHRONE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

COMMIT SYNCHRONE (commitSync) :
-> Bloque jusqu'à confirmation du broker
-> Garantie forte
-> Latence

consumer.commitSync();

COMMIT ASYNCHRONE (commitAsync) :
-> Non bloquant
-> Pas de garantie
-> Performance

consumer.commitAsync((offsets, exception) -> {
    if (exception != null) {
        System.err.println("Erreur commit: " + exception.getMessage());
    }
});


STRATÉGIES DE COMMIT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. COMMIT APRÈS CHAQUE MESSAGE (At-Most-Once)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
for (ConsumerRecord<String, String> record : records) {
    consumer.commitSync();  // Commit AVANT traitement
    processMessage(record);
}

-> Garantie : At-Most-Once (au plus une fois)
-> Si crash pendant traitement : message perdu
-> Pas de duplicata


2. COMMIT APRÈS TRAITEMENT (At-Least-Once)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
for (ConsumerRecord<String, String> record : records) {
    processMessage(record);
    consumer.commitSync();  // Commit APRÈS traitement
}

-> Garantie : At-Least-Once (au moins une fois)
-> Si crash pendant commit : message retraité
-> Possibilité de duplicata (idempotence nécessaire)


3. COMMIT PAR BATCH
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

for (ConsumerRecord<String, String> record : records) {
    processMessage(record);
}

consumer.commitSync();  // Commit une fois tout le batch traité

-> Plus performant
-> En cas de crash : retraitement du batch entier


4. COMMIT PAR PARTITION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();

for (ConsumerRecord<String, String> record : records) {
    processMessage(record);
    
    TopicPartition partition = new TopicPartition(record.topic(), record.partition());
    offsets.put(partition, new OffsetAndMetadata(record.offset() + 1));
}

consumer.commitSync(offsets);  // Commit par partition

-> Granularité fine
-> Meilleur contrôle


SEEK - POSITIONNER L'OFFSET MANUELLEMENT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// Aller au début de toutes les partitions
consumer.seekToBeginning(consumer.assignment());

// Aller à la fin (nouveaux messages seulement)
consumer.seekToEnd(consumer.assignment());

// Aller à un offset spécifique
TopicPartition partition = new TopicPartition("orders", 0);
consumer.seek(partition, 1250);

// Aller à un timestamp
Map<TopicPartition, Long> timestampsToSearch = new HashMap<>();
timestampsToSearch.put(partition, 1638316800000L);  // 1er Dec 2021
Map<TopicPartition, OffsetAndMetadata> offsets = consumer.offsetsForTimes(timestampsToSearch);
offsets.forEach((tp, offset) -> consumer.seek(tp, offset.offset()));


PATTERNS DE CONSOMMATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. CONSUMER BASIQUE (Single-threaded)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        processMessage(record);  // Traitement séquentiel
    }
    
    consumer.commitSync();
}

AVANTAGES :
[OK] Simple
[OK] Ordre garanti

INCONVÉNIENTS :
[X] Lent si traitement lourd
[X] Un seul thread


2. CONSUMER AVEC THREAD POOL (Multi-threaded processing)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ExecutorService executor = Executors.newFixedThreadPool(10);

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    if (records.isEmpty()) continue;
    
    CountDownLatch latch = new CountDownLatch(records.count());
    
    for (ConsumerRecord<String, String> record : records) {
        executor.submit(() -> {
            try {
                processMessage(record);
            } finally {
                latch.countDown();
            }
        });
    }
    
    latch.await();  // Attendre que tous les threads finissent
    consumer.commitSync();
}

AVANTAGES :
[OK] Traitement parallèle
[OK] Débit élevé

INCONVÉNIENTS :
[X] Ordre NON garanti
[X] Complexité accrue


3. CONSUMER AVEC PAUSE/RESUME
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        if (shouldBackpressure()) {
            // Pause la partition pour ralentir
            consumer.pause(Collections.singleton(
                new TopicPartition(record.topic(), record.partition())
            ));
        }
        
        processMessage(record);
        
        if (canResume()) {
            // Reprendre la consommation
            consumer.resume(consumer.paused());
        }
    }
    
    consumer.commitSync();
}


GESTION D'ERREURS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ERREURS RÉCUPÉRABLES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

try {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        try {
            processMessage(record);
            consumer.commitSync();
            
        } catch (RecoverableException e) {
            // Erreur temporaire : ne pas committer
            // Le message sera retraité
            log.warn("Erreur récupérable: {}", e.getMessage());
            
        } catch (UnrecoverableException e) {
            // Erreur permanente : envoyer en DLQ
            sendToDeadLetterQueue(record);
            consumer.commitSync();  // Commit pour passer au suivant
        }
    }
    
} catch (WakeupException e) {
    // Shutdown graceful
    log.info("Consumer en cours d'arrêt");
}


DEAD LETTER QUEUE (DLQ) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

private void sendToDeadLetterQueue(ConsumerRecord<String, String> record) {
    ProducerRecord<String, String> dlqRecord = new ProducerRecord<>(
        "orders-dlq",  // Topic DLQ
        record.key(),
        record.value()
    );
    
    // Ajouter des headers avec info d'erreur
    dlqRecord.headers()
        .add("original-topic", record.topic().getBytes())
        .add("original-partition", String.valueOf(record.partition()).getBytes())
        .add("original-offset", String.valueOf(record.offset()).getBytes())
        .add("error-reason", errorReason.getBytes());
    
    dlqProducer.send(dlqRecord);
}


SHUTDOWN GRACEFUL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

final KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

// Hook pour arrêt propre
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    System.out.println("Arrêt du consumer...");
    consumer.wakeup();  // Interrompt poll()
}));

try {
    consumer.subscribe(Collections.singletonList("orders"));
    
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        
        for (ConsumerRecord<String, String> record : records) {
            processMessage(record);
        }
        
        consumer.commitSync();
    }
    
} catch (WakeupException e) {
    // Attendre, wakeup() a été appelé
    System.out.println("Consumer interrompu proprement");
    
} finally {
    consumer.close();
    System.out.println("Consumer fermé [OK]");
}


BONNES PRATIQUES CONSUMER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] TOUJOURS utiliser enable.auto.commit=false
[OK] Commit APRÈS traitement réussi
[OK] Utiliser try-catch pour erreurs
[OK] Implémenter Dead Letter Queue
[OK] Surveiller le lag (décalage)
[OK] Configurer max.poll.interval.ms correctement
[OK] Implémenter shutdown graceful
[OK] Rendre le traitement IDEMPOTENT
[OK] Logger les offsets pour debug
[OK] Monitorer les rebalancing

[X] Ne PAS partager un consumer entre threads
[X] Ne PAS oublier de commit
[X] Ne PAS bloquer dans poll()
[X] Ne PAS traiter trop de messages entre deux poll()
[X] Ne PAS ignorer les WakeupException
"""


# [OK] PARTIE 6 : CONSUMER GROUPS ET PARALLÉLISME

"""
┌────────────────────────────────────────────────────────────────────────┐
│              CONSUMER GROUPS : PARALLÉLISME ET SCALABILITÉ             │
└────────────────────────────────────────────────────────────────────────┘

QU'EST-CE QU'UN CONSUMER GROUP ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un CONSUMER GROUP = ensemble de consumers partageant le même group.id

PRINCIPE FONDAMENTAL :
-> Chaque PARTITION est assignée à UN SEUL consumer du groupe
-> Un consumer peut avoir PLUSIEURS partitions
-> Permet le PARALLÉLISME et la SCALABILITÉ


EXEMPLE VISUEL :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic "orders" avec 3 partitions :
┌──────────────────────────────────────────────────────────────────┐
│ Partition 0: [Msg0][Msg3][Msg6]...                              │
│ Partition 1: [Msg1][Msg4][Msg7]...                              │
│ Partition 2: [Msg2][Msg5][Msg8]...                              │
└──────────────────────────────────────────────────────────────────┘

Consumer Group "analytics" avec 3 consumers :
┌──────────────────────────────────────────────────────────────────┐
│                                                                  │
│  Consumer A          Consumer B          Consumer C             │
│  (P0 assignée)       (P1 assignée)       (P2 assignée)          │
│      │                   │                   │                  │
│      [BLACK_DOWN-POINTING_TRIANGLE]                   [BLACK_DOWN-POINTING_TRIANGLE]                   [BLACK_DOWN-POINTING_TRIANGLE]                  │
│  [Msg0, Msg3...]     [Msg1, Msg4...]     [Msg2, Msg5...]       │
└──────────────────────────────────────────────────────────────────┘

-> Chaque consumer lit UNE partition
-> Parallélisme maximum = 3 (nombre de partitions)


CAS 1 : UN SEUL CONSUMER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Consumer Group "analytics" avec 1 consumer :
┌──────────────────────────────────────────────────────────────────┐
│                        Consumer A                                │
│                  (P0, P1, P2 assignées)                          │
│      │                │                │                         │
│      [BLACK_DOWN-POINTING_TRIANGLE]                [BLACK_DOWN-POINTING_TRIANGLE]                [BLACK_DOWN-POINTING_TRIANGLE]                         │
│  [Msg0, Msg3...]  [Msg1, Msg4...]  [Msg2, Msg5...]             │
└──────────────────────────────────────────────────────────────────┘

-> Un consumer lit TOUTES les partitions
-> Pas de parallélisme
-> Traitement séquentiel


CAS 2 : DEUX CONSUMERS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Consumer Group "analytics" avec 2 consumers :
┌──────────────────────────────────────────────────────────────────┐
│           Consumer A                    Consumer B               │
│       (P0, P1 assignées)               (P2 assignée)             │
│      │        │                             │                    │
│      [BLACK_DOWN-POINTING_TRIANGLE]        [BLACK_DOWN-POINTING_TRIANGLE]                             [BLACK_DOWN-POINTING_TRIANGLE]                    │
│  [Msg0...]  [Msg1...]                   [Msg2...]                │
└──────────────────────────────────────────────────────────────────┘

-> Répartition inégale (A : 2 partitions, B : 1 partition)
-> A traite plus de messages que B


CAS 3 : PLUS DE CONSUMERS QUE DE PARTITIONS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic avec 3 partitions, Consumer Group avec 5 consumers :
┌──────────────────────────────────────────────────────────────────────┐
│  Consumer A     Consumer B     Consumer C     Consumer D  Consumer E │
│  (P0)           (P1)           (P2)           (IDLE)      (IDLE)     │
│    │              │              │                                   │
│    [BLACK_DOWN-POINTING_TRIANGLE]              [BLACK_DOWN-POINTING_TRIANGLE]              [BLACK_DOWN-POINTING_TRIANGLE]                                   │
│ [Msg0...]      [Msg1...]      [Msg2...]                              │
└──────────────────────────────────────────────────────────────────────┘

-> 2 consumers INUTILISÉS (idle)
-> Maximum de consumers utiles = nombre de partitions

[ATTENTION]  IMPORTANT :
Nombre optimal de consumers = nombre de partitions
Pas de parallélisme supplémentaire au-delà !


MULTIPLES CONSUMER GROUPS (Indépendants)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic "orders" lu par PLUSIEURS groups :

Consumer Group "analytics":
┌──────────────────────────────────────────────────────────────────┐
│  Consumer A-1      Consumer A-2      Consumer A-3                │
│  (P0)              (P1)              (P2)                        │
│    │                 │                 │                         │
│    [BLACK_DOWN-POINTING_TRIANGLE]                 [BLACK_DOWN-POINTING_TRIANGLE]                 [BLACK_DOWN-POINTING_TRIANGLE]                         │
│ [Msg0...]         [Msg1...]         [Msg2...]                    │
└──────────────────────────────────────────────────────────────────┘
          v Traite pour analytics (dashboards, reports)


Consumer Group "billing":
┌──────────────────────────────────────────────────────────────────┐
│  Consumer B-1      Consumer B-2      Consumer B-3                │
│  (P0)              (P1)              (P2)                        │
│    │                 │                 │                         │
│    [BLACK_DOWN-POINTING_TRIANGLE]                 [BLACK_DOWN-POINTING_TRIANGLE]                 [BLACK_DOWN-POINTING_TRIANGLE]                         │
│ [Msg0...]         [Msg1...]         [Msg2...]                    │
└──────────────────────────────────────────────────────────────────┘
          v Traite pour billing (factures, paiements)


Consumer Group "notifications":
┌──────────────────────────────────────────────────────────────────┐
│                      Consumer C-1                                │
│                   (P0, P1, P2)                                   │
│      │                │                │                         │
│      [BLACK_DOWN-POINTING_TRIANGLE]                [BLACK_DOWN-POINTING_TRIANGLE]                [BLACK_DOWN-POINTING_TRIANGLE]                         │
│  [Msg0...]         [Msg1...]         [Msg2...]                   │
└──────────────────────────────────────────────────────────────────┘
          v Traite pour notifications (emails, SMS)

-> Chaque group lit TOUS les messages INDÉPENDAMMENT
-> Offsets séparés par group
-> Un group n'affecte pas les autres


REBALANCING (Réassignation des partitions)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉCLENCHEURS DE REBALANCING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Nouveau consumer rejoint le groupe
2. Consumer quitte le groupe (crash ou shutdown)
3. Consumer inactif (pas de heartbeat pendant session.timeout.ms)
4. Ajout/suppression de partitions au topic
5. Consumer ne poll() pas pendant max.poll.interval.ms


PROCESSUS DE REBALANCING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AVANT (état stable) :
┌────────────────────────────────────────────────────────────────┐
│  Consumer A       Consumer B       Consumer C                  │
│  (P0)             (P1)             (P2)                        │
└────────────────────────────────────────────────────────────────┘

Consumer B CRASH [IMPACT]

PENDANT LE REBALANCING :
┌────────────────────────────────────────────────────────────────┐
│  Consumer A       [IMPACT] CRASH         Consumer C                  │
│  (Paused)                          (Paused)                    │
│                                                                │
│  -> Coordinator détecte la panne                                │
│  -> Déclenche le rebalancing                                    │
│  -> Tous les consumers ARRÊTENT de consommer                    │
└────────────────────────────────────────────────────────────────┘

APRÈS (nouvel état stable) :
┌────────────────────────────────────────────────────────────────┐
│  Consumer A                       Consumer C                   │
│  (P0, P1)                         (P2)                         │
│                                                                │
│  -> Partitions réassignées                                      │
│  -> Consommation reprend                                        │
└────────────────────────────────────────────────────────────────┘

DURÉE :
-> Quelques secondes généralement
-> Dépend de la taille du groupe et du nombre de partitions

IMPACT :
[X] STOP THE WORLD : tous les consumers arrêtent pendant le rebalancing
[X] Perte temporaire de débit
[X] Possibilité de retraiter des messages (si pas de commit avant rebalancing)


STRATÉGIES DE RÉASSIGNATION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. RANGE ASSIGNOR (défaut avant Kafka 3.0)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Assigne des plages de partitions
-> Peut être déséquilibré

Exemple : Topic avec 5 partitions, 3 consumers
Consumer A : P0, P1
Consumer B : P2, P3
Consumer C : P4

2. ROUND ROBIN ASSIGNOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Distribution circulaire
-> Plus équilibré

Exemple : Topic avec 5 partitions, 3 consumers
Consumer A : P0, P3
Consumer B : P1, P4
Consumer C : P2

3. STICKY ASSIGNOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Garde les assignations existantes autant que possible
-> Réduit les mouvements de partitions
-> Rebalancing plus rapide

4. COOPERATIVE STICKY (défaut depuis Kafka 3.0)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Réassigne SEULEMENT les partitions nécessaires
-> Pas de stop-the-world
-> Consumers continuent sur partitions non affectées

CONFIGURATION :
props.put("partition.assignment.strategy", 
          "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");


INCREMENTAL COOPERATIVE REBALANCING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AVANT (Eager Rebalancing) :
1. Tous les consumers RELÂCHENT toutes leurs partitions
2. Pause complète
3. Réassignation
4. Reprise

APRÈS (Cooperative Rebalancing) :
1. Seules les partitions à déplacer sont relâchées
2. Autres partitions CONTINUENT à être consommées
3. Réassignation incrémentale
4. Impact minimal

AVANTAGES :
[OK] Pas de stop-the-world complet
[OK] Débit maintenu
[OK] Latence réduite


LISTENERS DE REBALANCING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Pour exécuter du code avant/après rebalancing :

consumer.subscribe(
    Collections.singletonList("orders"),
    new ConsumerRebalanceListener() {
        
        @Override
        public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
            // Appelé AVANT de perdre des partitions
            System.out.println("Partitions révoquées: " + partitions);
            
            // Commit les offsets avant de perdre les partitions
            consumer.commitSync();
            
            // Sauvegarder l'état si nécessaire
            saveState();
        }
        
        @Override
        public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
            // Appelé APRÈS avoir reçu de nouvelles partitions
            System.out.println("Partitions assignées: " + partitions);
            
            // Restaurer l'état si nécessaire
            restoreState();
        }
    }
);


MONITORING DU LAG
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

LAG = Nombre de messages en retard (non consommés)

┌────────────────────────────────────────────────────────────────┐
│ Partition 0                                                    │
│ [Msg0][Msg1][Msg2][Msg3][Msg4][Msg5][Msg6][Msg7][Msg8][Msg9]   │
│                         [BLACK_UP-POINTING_TRIANGLE]                            [BLACK_UP-POINTING_TRIANGLE]         │
│                         │                            │         │
│                   Current Offset (3)        Log End Offset (10)│
│                                                                │
│                   LAG = 10 - 3 = 7 messages                    │
└────────────────────────────────────────────────────────────────┘

COMMANDE POUR VOIR LE LAG :
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe \
  --group my-consumer-group

Output:
GROUP           TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
my-group        orders    0          1250            1250            0
my-group        orders    1          980             1100            120  [ATTENTION]
my-group        orders    2          1500            1500            0

-> Partition 1 a un lag de 120 messages


CAUSES DU LAG :
[X] Consumer trop lent (traitement lourd)
[X] Pas assez de consumers (sous-provisionné)
[X] Producer trop rapide
[X] Consumer en erreur/crash répétés
[X] Rebalancing fréquents

SOLUTIONS :
[OK] Augmenter le nombre de partitions
[OK] Ajouter des consumers au groupe
[OK] Optimiser le traitement
[OK] Utiliser le multi-threading
[OK] Scaler horizontalement


COORDINATION ET GROUP COORDINATOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque consumer group a un GROUP COORDINATOR :
-> Un broker Kafka désigné
-> Gère les heartbeats
-> Orchestre les rebalancings
-> Stocke les offsets committés

TOPIC INTERNE : __consumer_offsets
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Topic spécial Kafka (50 partitions par défaut)
-> Stocke les offsets de TOUS les consumer groups
-> Compacté (log compaction)

Format d'une entrée :
Key: (group.id, topic, partition)
Value: (offset, metadata, timestamp)


EXEMPLE COMPLET : DÉPLOIEMENT D'UN CONSUMER GROUP
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SCÉNARIO : E-commerce avec topic "orders" (12 partitions)

DÉPLOIEMENT :
┌────────────────────────────────────────────────────────────────┐
│                   CONSUMER GROUP "order-processor"             │
│                                                                │
│  Instance 1 (Container)     Instance 2 (Container)             │
│  Consumer A                  Consumer B                        │
│  Partitions: 0,1,2,3,4,5    Partitions: 6,7,8,9,10,11          │
│                                                                │
│  -> Débit : 50K msgs/sec      -> Débit : 50K msgs/sec            │
│  -> Répartition équitable     -> Total : 100K msgs/sec           │
└────────────────────────────────────────────────────────────────┘

SI BESOIN DE PLUS DE DÉBIT :
-> Ajouter Instance 3, 4, ... jusqu'à 12 (nombre de partitions)

SI INSTANCE 2 CRASH :
-> Rebalancing automatique
-> Instance 1 prend toutes les 12 partitions temporairement
-> Débit réduit de moitié jusqu'à récupération


BONNES PRATIQUES CONSUMER GROUPS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Nombre de consumers = nombre de partitions (optimum)
[OK] Utiliser Cooperative Sticky Assignor
[OK] Monitorer le lag constamment
[OK] Implémenter des listeners de rebalancing
[OK] Commit dans onPartitionsRevoked
[OK] Rendre le traitement idempotent
[OK] Configurer correctement session.timeout.ms et max.poll.interval.ms
[OK] Un group.id unique par application/service
[OK] Tester les scénarios de crash
[OK] Éviter les rebalancements fréquents

[X] Ne PAS avoir trop de consumers (idle)
[X] Ne PAS partager le group.id entre applications différentes
[X] Ne PAS oublier de commit avant rebalancing
[X] Ne PAS bloquer trop longtemps le thread de poll
[X] Ne PAS traiter trop lentement (cause du lag)
"""


# [OK] PARTIE 7 : KAFKA SOUS LE CAPOT (LOG, SEGMENTS, INDEX)

"""
┌────────────────────────────────────────────────────────────────────────┐
│              KAFKA INTERNALS : COMMENT ÇA MARCHE VRAIMENT ?            │
└────────────────────────────────────────────────────────────────────────┘

KAFKA = LOG DISTRIBUÉ IMMUTABLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION D'UN LOG :
-> Fichier append-only (ajout seulement à la fin)
-> Ordonné par ordre d'arrivée
-> Immuable (une fois écrit, ne change jamais)


STRUCTURE SUR DISQUE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

HIÉRARCHIE :
Kafka Data Directory (/var/kafka-logs/)
  └─ Topic-Partition directories
      └─ Segment files (.log, .index, .timeindex)


EXEMPLE RÉEL :
/var/kafka-logs/
├── orders-0/                         <- Partition 0 du topic "orders"
│   ├── 00000000000000000000.log      <- Segment 1 (messages)
│   ├── 00000000000000000000.index    <- Index offset -> position
│   ├── 00000000000000000000.timeindex<- Index timestamp -> offset
│   ├── 00000000000000050000.log      <- Segment 2
│   ├── 00000000000000050000.index
│   ├── 00000000000000050000.timeindex
│   ├── 00000000000000100000.log      <- Segment 3
│   ├── 00000000000000100000.index
│   ├── 00000000000000100000.timeindex
│   ├── leader-epoch-checkpoint       <- Info leader epoch
│   └── partition.metadata            <- Métadonnées partition
│
├── orders-1/                         <- Partition 1
│   ├── ...
│
├── orders-2/                         <- Partition 2
│   ├── ...
│
├── __consumer_offsets-0/             <- Topic interne des offsets
│   ├── ...
│
└── recovery-point-offset-checkpoint   <- Point de récupération
└── replication-offset-checkpoint      <- Point de réplication


SEGMENTS - DÉCOUPAGE DU LOG
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un LOG de partition est divisé en SEGMENTS :

Partition 0 du topic "orders" :
┌─────────────────────────────────────────────────────────────────┐
│ SEGMENT 1             SEGMENT 2             SEGMENT 3 (actif)   │
│ [0-49999]             [50000-99999]         [100000-...]        │
│                                                                 │
│ 00000...00000.log     00000...50000.log     00000...100000.log │
│ (FERMÉ)               (FERMÉ)               (ACTIF, écriture)  │
└─────────────────────────────────────────────────────────────────┘

RÈGLES DE ROTATION (création d'un nouveau segment) :
1. Taille atteinte (log.segment.bytes = 1 GB par défaut)
2. Temps écoulé (log.roll.ms ou log.roll.hours = 7 jours par défaut)

Configuration :
log.segment.bytes=1073741824     # 1 GB
log.roll.hours=168               # 7 jours
log.roll.ms=604800000            # 7 jours en ms


POURQUOI DES SEGMENTS ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. SUPPRESSION EFFICACE (Rétention) :
   -> Supprimer un vieux segment = supprimer un fichier
   -> Pas besoin de réécrire tout le log

2. COMPACTION (Log Compaction) :
   -> Compacter par segment
   -> Garder seulement la dernière valeur par clé

3. PERFORMANCE :
   -> Fichiers de taille raisonnable
   -> OS peut cacher efficacement
   -> Flush séquentiel


FORMAT D'UN MESSAGE DANS LE LOG (.log file)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque message stocké contient :

┌─────────────────────────────────────────────────────────────────┐
│ MESSAGE (Record Batch)                                          │
│                                                                 │
│ ┌─────────────────────┐                                         │
│ │ BASE OFFSET         │ 8 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ LENGTH              │ 4 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ PARTITION LEADER    │ 4 bytes                                 │
│ │ EPOCH               │                                         │
│ ├─────────────────────┤                                         │
│ │ MAGIC               │ 1 byte (version format)                 │
│ ├─────────────────────┤                                         │
│ │ CRC                 │ 4 bytes (checksum)                      │
│ ├─────────────────────┤                                         │
│ │ ATTRIBUTES          │ 2 bytes (compression, timestamp type)   │
│ ├─────────────────────┤                                         │
│ │ LAST OFFSET DELTA   │ 4 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ FIRST TIMESTAMP     │ 8 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ MAX TIMESTAMP       │ 8 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ PRODUCER ID         │ 8 bytes (idempotence)                   │
│ ├─────────────────────┤                                         │
│ │ PRODUCER EPOCH      │ 2 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ BASE SEQUENCE       │ 4 bytes (idempotence)                   │
│ ├─────────────────────┤                                         │
│ │ RECORDS COUNT       │ 4 bytes                                 │
│ ├─────────────────────┤                                         │
│ │ RECORDS             │ Variable (messages compressés si config)│
│ │  ├─ Key             │                                         │
│ │  ├─ Value           │                                         │
│ │  ├─ Headers         │                                         │
│ │  └─ Timestamp       │                                         │
│ └─────────────────────┘                                         │
└─────────────────────────────────────────────────────────────────┘

-> Batch de messages pour efficacité
-> Compression au niveau du batch
-> Checksum CRC pour intégrité


INDEX (.index file)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

L'INDEX permet de trouver rapidement un message par son offset

STRUCTURE (index offset -> position) :
┌─────────────────────────────────────────────────────────────────┐
│ 00000000000000000000.index                                      │
│                                                                 │
│ Offset    Physical Position (in .log file)                     │
│ ──────    ────────────────────────────────────────              │
│ 0         0                                                     │
│ 1000      102400                                                │
│ 2000      204800                                                │
│ 3000      307200                                                │
│ ...                                                             │
└─────────────────────────────────────────────────────────────────┘

-> SPARSE INDEX : pas tous les offsets, seulement échantillons
-> Intervalle défini par log.index.interval.bytes (4 KB par défaut)


RECHERCHE D'UN MESSAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Exemple : Chercher message avec offset 2500

1. Identifier le segment qui contient offset 2500
   -> Segment 00000000000000000000.log (contient 0-49999)

2. Lire l'index 00000000000000000000.index
   -> Recherche binaire (O(log n))
   -> Trouver l'entrée la plus proche : offset 2000 -> position 204800

3. Lire le .log file à partir de position 204800
   -> Lire séquentiellement jusqu'à offset 2500

4. Retourner le message

PERFORMANCE :
-> Recherche binaire dans index : très rapide
-> Lecture séquentielle dans log : optimisée par OS cache


TIME INDEX (.timeindex file)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Permet de chercher par TIMESTAMP

STRUCTURE (timestamp -> offset) :
┌─────────────────────────────────────────────────────────────────┐
│ 00000000000000000000.timeindex                                  │
│                                                                 │
│ Timestamp                  Offset                               │
│ ─────────────────          ──────                               │
│ 1638316800000 (Dec 1)      0                                    │
│ 1638403200000 (Dec 2)      1000                                 │
│ 1638489600000 (Dec 3)      2000                                 │
│ ...                                                             │
└─────────────────────────────────────────────────────────────────┘

USE CASE :
-> Chercher messages depuis une date
-> consumer.seek() avec timestamp


ZERO-COPY - PERFORMANCE RÉSEAU
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TRADITIONNEL (avec copy) :
┌─────────┐     ┌─────────┐     ┌─────────┐     ┌─────────┐
│  DISK   │ ──-> │ OS Cache│ ──-> │App (JVM)│ ──-> │ Network │
└─────────┘     └─────────┘     └─────────┘     └─────────┘
                    v                 v
                  Copy 1            Copy 2

-> 2 copies en mémoire (overhead)

KAFKA ZERO-COPY (sendfile) :
┌─────────┐     ┌─────────┐     ┌─────────┐
│  DISK   │ ──-> │ OS Cache│ ──-> │ Network │
└─────────┘     └─────────┘     └─────────┘

-> 0 copy en userspace !
-> Données vont directement du kernel au socket réseau
-> Ultra performant

COMMENT ?
-> Java NIO : FileChannel.transferTo()
-> Linux : sendfile() system call


PAGE CACHE - UTILISATION INTELLIGENTE DE LA RAM
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Kafka utilise le PAGE CACHE de l'OS :
-> Pas de cache applicatif (pas de gestion JVM)
-> Laisse l'OS gérer la mémoire
-> Lecture/écriture très rapide si données en cache

AVANTAGES :
[OK] Survit aux redémarrages de Kafka (cache OS persiste)
[OK] Pas de GC (Garbage Collection) JVM
[OK] Mémoire libérée automatiquement sous pression
[OK] Warm up automatique après redémarrage


ÉCRITURE SÉQUENTIELLE - POURQUOI KAFKA EST RAPIDE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ÉCRITURE ALÉATOIRE (Base de données classique) :
-> Seek + Write à différents endroits du disque
-> Lent sur HDD (seek time ~10ms)
-> Débit : ~100 ops/sec

ÉCRITURE SÉQUENTIELLE (Kafka) :
-> Append seulement à la fin du fichier
-> Pas de seek
-> Débit : ~600 MB/sec sur HDD, 1+ GB/sec sur SSD

DISQUE SÉQUENTIEL > RAM ALÉATOIRE !
-> Oui, écrire séquentiellement sur disque peut être plus rapide
   que lire aléatoirement en RAM (CPU cache misses)


LOG COMPACTION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DEUX POLITIQUES DE RÉTENTION :

1. TIME/SIZE-BASED DELETION (défaut) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Supprimer segments après X heures/jours
-> Ou supprimer si taille totale > Y

log.retention.hours=168        # 7 jours
log.retention.bytes=1073741824 # 1 GB

2. LOG COMPACTION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Garder seulement la DERNIÈRE valeur de chaque clé
-> Utile pour changelog, state store

EXEMPLE :
AVANT COMPACTION :
Key=user-1, Value={"name": "Jean", "age": 30}      Offset 0
Key=user-2, Value={"name": "Marie", "age": 25}     Offset 1
Key=user-1, Value={"name": "Jean", "age": 31}      Offset 2  <- Update
Key=user-3, Value={"name": "Paul", "age": 40}      Offset 3
Key=user-2, Value=null                              Offset 4  <- Tombstone (delete)

APRÈS COMPACTION :
Key=user-1, Value={"name": "Jean", "age": 31}      Offset 2  <- Dernière valeur
Key=user-3, Value={"name": "Paul", "age": 40}      Offset 3
# user-2 supprimé (tombstone)

CONFIGURATION :
log.cleanup.policy=compact

USE CASES :
-> CDC (Change Data Capture)
-> Kafka Streams state stores
-> Materialisation de vues


FLUSH POLICY - QUAND ÉCRIRE SUR DISQUE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Par défaut, Kafka ne flush PAS immédiatement :
-> Écritures bufferisées en RAM (page cache)
-> OS décide quand écrire sur disque

CONFIGURATION :
log.flush.interval.messages=10000  # Flush après 10K messages
log.flush.interval.ms=1000         # Flush toutes les secondes

RECOMMANDATION :
-> NE PAS configurer flush explicitement
-> Laisser l'OS gérer (meilleure performance)
-> Utiliser replication.factor >= 3 pour durabilité


RÉSUMÉ PERFORMANCE KAFKA :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Écriture séquentielle (append-only)
[OK] Zero-copy (sendfile)
[OK] Page cache OS (pas de cache applicatif)
[OK] Batching (messages groupés)
[OK] Compression (moins de bytes)
[OK] Partitioning (parallélisme)
[OK] Pas de suppression en place (segments entiers)

RÉSULTAT :
-> Millions de messages par seconde
-> Latence : quelques millisecondes
-> Débit : plusieurs GB/sec par broker
"""

Je continue le guide Kafka avec les parties restantes :

```txt
# [OK] PARTIE 8 : RÉPLICATION ET HAUTE DISPONIBILITÉ

"""
┌────────────────────────────────────────────────────────────────────────┐
│              RÉPLICATION : HAUTE DISPONIBILITÉ ET DURABILITÉ           │
└────────────────────────────────────────────────────────────────────────┘

POURQUOI LA RÉPLICATION ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SANS RÉPLICATION :
┌────────────────────────────────────────────────────────────────┐
│  BROKER 1 (seul)                                               │
│  orders-0 (UNIQUE COPY)                                        │
│  [Msg0][Msg1][Msg2]...                                         │
└────────────────────────────────────────────────────────────────┘
              │
              [BLACK_DOWN-POINTING_TRIANGLE]
           [IMPACT] CRASH
              │
              [BLACK_DOWN-POINTING_TRIANGLE]
        DONNÉES PERDUES [X]

AVEC RÉPLICATION (factor = 3) :
┌──────────────────────────────────────────────────────────────────┐
│  BROKER 1 (Leader)    BROKER 2 (Follower)    BROKER 3 (Follower)│
│  orders-0             orders-0               orders-0            │
│  [Msg0][Msg1]...      [Msg0][Msg1]...        [Msg0][Msg1]...    │
└──────────────────────────────────────────────────────────────────┘
              │
              [BLACK_DOWN-POINTING_TRIANGLE]
           [IMPACT] CRASH
              │
              [BLACK_DOWN-POINTING_TRIANGLE]
┌──────────────────────────────────────────────────────────────────┐
│  BROKER 2 devient Leader [OK]    BROKER 3 (Follower)              │
│  orders-0                      orders-0                          │
│  [Msg0][Msg1]...               [Msg0][Msg1]...                   │
└──────────────────────────────────────────────────────────────────┘
         AUCUNE PERTE DE DONNÉES [OK]


REPLICATION FACTOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Nombre de copies d'une partition (leader + followers)

EXEMPLE : Topic "orders" avec 3 partitions, replication.factor=3

┌─────────────────────────────────────────────────────────────────┐
│                        KAFKA CLUSTER                            │
│                                                                 │
│  BROKER 1           BROKER 2           BROKER 3                │
│  ─────────          ─────────          ─────────               │
│  P0 (Leader)        P0 (Follower)      P0 (Follower)           │
│  P1 (Follower)      P1 (Leader)        P1 (Follower)           │
│  P2 (Follower)      P2 (Follower)      P2 (Leader)             │
└─────────────────────────────────────────────────────────────────┘

-> Chaque partition a 3 copies
-> 1 Leader, 2 Followers
-> Leaders répartis équitablement


CRÉATION AVEC REPLICATION :
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create \
  --topic orders \
  --partitions 3 \
  --replication-factor 3

Configuration broker par défaut :
default.replication.factor=3


RECOMMANDATIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Production : replication.factor = 3 (minimum)
[OK] Critique : replication.factor = 5
[OK] Développement : replication.factor = 1 (acceptable)

[ATTENTION]  IMPORTANT :
replication.factor NE PEUT PAS dépasser le nombre de brokers
-> 3 brokers max = replication.factor max de 3


LEADER ET FOLLOWERS - DÉTAIL DU FONCTIONNEMENT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLES :

LEADER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Reçoit TOUTES les écritures (producers)
-> Sert TOUTES les lectures (consumers, par défaut)
-> Coordonne la réplication vers followers
-> Gère les High Water Mark et Log End Offset

FOLLOWERS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Répliquent les données depuis le leader
-> PULL continuellement (fetch requests)
-> NE servent PAS les lectures (par défaut)
-> Candidats pour devenir leader


FLUX DE RÉPLICATION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Producer envoie un message
   v
2. Leader reçoit et écrit dans son log local
   v
3. Followers envoient fetch requests au leader
   v
4. Leader répond avec nouveaux messages
   v
5. Followers écrivent dans leurs logs locaux
   v
6. Followers envoient ACK au leader
   v
7. Leader met à jour High Water Mark
   v
8. Leader répond au producer (si acks=all)

┌──────────────────────────────────────────────────────────────────┐
│                                                                  │
│  PRODUCER                                                        │
│      │                                                           │
│      │ (1) Send message                                         │
│      [BLACK_DOWN-POINTING_TRIANGLE]                                                           │
│  ┌─────────────────┐                                            │
│  │  LEADER (B1)    │                                            │
│  │  P0             │                                            │
│  │  LEO: 105 ──────┼───────────────────┐                        │
│  │  HW:  103       │                   │                        │
│  └─────────────────┘                   │                        │
│         │                              │                        │
│         │ (3) Fetch                    │ (3) Fetch              │
│         [BLACK_DOWN-POINTING_TRIANGLE]                              [BLACK_DOWN-POINTING_TRIANGLE]                        │
│  ┌─────────────┐              ┌─────────────┐                  │
│  │FOLLOWER (B2)│              │FOLLOWER (B3)│                  │
│  │  P0         │              │  P0         │                  │
│  │  LEO: 104   │              │  LEO: 103   │                  │
│  └─────────────┘              └─────────────┘                  │
│         │                              │                        │
│         └──────────────────────────────┘                        │
│                       │                                          │
│                       │ (6) ACK when caught up                  │
│                       [BLACK_DOWN-POINTING_TRIANGLE]                                          │
│                  HW updated to 104                               │
└──────────────────────────────────────────────────────────────────┘

LEO (Log End Offset) : Dernier offset écrit dans le log
HW (High Water Mark) : Dernier offset répliqué sur tous les ISR


IN-SYNC REPLICAS (ISR)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ISR = Liste des répliques SYNCHRONISÉES avec le leader

CRITÈRES POUR ÊTRE DANS L'ISR :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Replica doit être ALIVE (heartbeat récent)
2. Replica ne doit pas avoir trop de LAG (replica.lag.time.max.ms)

Configuration broker :
replica.lag.time.max.ms=30000  # 30 secondes

Si un follower prend du retard > 30s :
-> RETIRÉ de l'ISR
-> Ne participe plus aux ACKs
-> Ne peut pas devenir leader


EXEMPLE D'ISR :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

État initial (tous synchronisés) :
Partition orders-0
Leader: Broker 1
ISR: [Broker 1, Broker 2, Broker 3]  <- Tous en sync

Broker 3 devient lent (network issue) :
Partition orders-0
Leader: Broker 1
ISR: [Broker 1, Broker 2]  <- Broker 3 retiré

Broker 3 rattrape son retard :
Partition orders-0
Leader: Broker 1
ISR: [Broker 1, Broker 2, Broker 3]  <- Broker 3 réintégré


MIN IN-SYNC REPLICAS (min.insync.replicas)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Nombre MINIMUM de répliques ISR pour accepter une écriture

Configuration (niveau topic ou broker) :
min.insync.replicas=2

COMPORTEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SCÉNARIO 1 : replication.factor=3, min.insync.replicas=2
ISR: [B1, B2, B3]  -> [OK] OK (3 >= 2)

SCÉNARIO 2 : B3 tombe
ISR: [B1, B2]  -> [OK] OK (2 >= 2)

SCÉNARIO 3 : B2 tombe aussi
ISR: [B1]  -> [X] ERREUR NotEnoughReplicasException
            -> Écritures REFUSÉES
            -> Lectures OK

GARANTIE :
Avec acks=all + min.insync.replicas=2 :
-> Au moins 2 copies avant de confirmer
-> Survit à la perte de (replication.factor - min.insync.replicas) brokers


CONFIGURATIONS RECOMMANDÉES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

HAUTE DISPONIBILITÉ (favoriser disponibilité) :
replication.factor=3
min.insync.replicas=1
acks=1

-> Accepte les écritures même si 2 brokers down
-> Risque de perte si leader crash avant réplication

HAUTE DURABILITÉ (favoriser durabilité) :
replication.factor=3
min.insync.replicas=2
acks=all

-> Garantie : 2 copies avant confirmation
-> Peut refuser écritures si < 2 ISR
-> Recommandé en PRODUCTION

ULTRA HAUTE DURABILITÉ (données critiques) :
replication.factor=5
min.insync.replicas=3
acks=all

-> Survit à la perte de 2 brokers
-> Coût : plus de stockage, latence


ÉLECTION DU LEADER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

QUAND UNE ÉLECTION SE DÉCLENCHE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Leader actuel crash
2. Leader actuel est shutdown gracefully
3. Broker leader devient isolé (network partition)
4. Changement de preferred leader (rebalancing)


PROCESSUS D'ÉLECTION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Controller (broker spécial) détecte la panne du leader
   v
2. Controller sélectionne un nouveau leader parmi l'ISR
   v
3. Critères de sélection :
   - Doit être dans l'ISR
   - Préférence : preferred leader (premier dans replica list)
   - Sinon : premier réplique dans l'ISR
   v
4. Controller met à jour ZooKeeper/KRaft
   v
5. Controller notifie tous les brokers
   v
6. Nouveau leader commence à servir les requêtes
   v
7. Followers se synchronisent avec le nouveau leader

DURÉE TYPIQUE : 2-5 secondes


UNCLEAN LEADER ELECTION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

QUE SE PASSE-T-IL SI TOUS LES ISR TOMBENT ?

SCÉNARIO :
Topic avec replication.factor=3
Leader : Broker 1
ISR : [Broker 1, Broker 2, Broker 3]

[IMPACT] Broker 1 crash
ISR : [Broker 2, Broker 3]
Leader : Broker 2 (élu)

[IMPACT] Broker 2 crash
ISR : [Broker 3]
Leader : Broker 3 (élu)

[IMPACT] Broker 3 crash
ISR : []  <- Aucun ISR disponible !

OPTIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

unclean.leader.election.enable=false (défaut, recommandé)
-> ATTENDRE qu'un ISR revienne
-> Partition INDISPONIBLE en écriture/lecture
-> Pas de perte de données [OK]

unclean.leader.election.enable=true
-> Élire un follower HORS ISR (out-of-sync)
-> Partition disponible [OK]
-> Risque de PERTE DE DONNÉES [X]

TRADE-OFF :
Disponibilité vs Durabilité


PREFERRED LEADER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Le "preferred leader" = premier réplique dans la liste des répliques

Partition orders-0
Replicas: [Broker 1, Broker 2, Broker 3]
Preferred Leader: Broker 1


POURQUOI C'EST IMPORTANT ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Kafka tente de RÉPARTIR les leaders équitablement :

ÉTAT INITIAL (équilibré) :
Broker 1: Leader de P0, P3
Broker 2: Leader de P1, P4
Broker 3: Leader de P2, P5

APRÈS PANNE DE BROKER 1 :
Broker 2: Leader de P0, P1, P3, P4  <- Déséquilibré !
Broker 3: Leader de P2, P5

QUAND BROKER 1 REVIENT :
-> P0 et P3 redeviennent leaders sur Broker 1 (preferred leader)
-> Rééquilibrage automatique

Configuration :
auto.leader.rebalance.enable=true  # Rééquilibrage auto (défaut)
leader.imbalance.check.interval.seconds=300  # Vérification toutes les 5 min


RACK AWARENESS (Data Center Awareness)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
Si tous les répliques sont dans le même rack/datacenter :
-> Panne du rack = perte de toutes les copies

SOLUTION : RACK AWARENESS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Configuration broker :
broker.rack=rack-1   # Broker 1 dans rack 1
broker.rack=rack-2   # Broker 2 dans rack 2
broker.rack=rack-3   # Broker 3 dans rack 3

COMPORTEMENT :
Kafka répartit les répliques sur DIFFÉRENTS racks

Partition orders-0 (replication.factor=3)
Replicas:
  - Broker 1 (rack-1)  <- Leader
  - Broker 2 (rack-2)  <- Follower
  - Broker 3 (rack-3)  <- Follower

-> Panne d'un rack = seulement 1 réplique perdue [OK]


FOLLOWER FETCHING (Lecture depuis follower)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DEPUIS KAFKA 2.4 : Possibilité de lire depuis FOLLOWER

POURQUOI ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Réduire latence (lire depuis le replica le plus proche)
-> Réduire charge réseau (inter-datacenter)
-> Réduire coût (éviter data transfer costs cloud)

EXEMPLE : Déploiement Multi-région

┌────────────────────────────────────────────────────────────────┐
│  AWS us-east-1                     AWS eu-west-1              │
│                                                               │
│  Broker 1 (Leader)                 Broker 3 (Follower)       │
│  orders-0                          orders-0                  │
│      [BLACK_UP-POINTING_TRIANGLE]                                  [BLACK_UP-POINTING_TRIANGLE]                    │
│      │                                  │                    │
│  Consumer US                        Consumer EU              │
│  (lit depuis leader)                (lit depuis follower)    │
│  Latence: 2ms                       Latence: 5ms             │
│                                     (évite 150ms inter-DC)   │
└────────────────────────────────────────────────────────────────┘

CONFIGURATION :
# Broker
broker.rack=us-east-1
replica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector

# Consumer
client.rack=us-east-1  # Consumer dans us-east-1 lit depuis us-east-1


RÉPLICATION MONITORING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MÉTRIQUES IMPORTANTES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Under-Replicated Partitions
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Partitions avec moins de répliques ISR que replication.factor

[ATTENTION]  ALERTE CRITIQUE si > 0

Vérification :
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe \
  --under-replicated-partitions

JMX Metric :
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions

2. Offline Partitions
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Partitions sans leader (indisponibles)

[ATTENTION]  ALERTE CRITIQUE si > 0

JMX Metric :
kafka.controller:type=KafkaController,name=OfflinePartitionsCount

3. ISR Shrink/Expand Rate
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Fréquence de changement de l'ISR

[ATTENTION]  ALERTE si trop fréquent (instabilité)

JMX Metrics :
kafka.server:type=ReplicaManager,name=IsrShrinksPerSec
kafka.server:type=ReplicaManager,name=IsrExpandsPerSec

4. Replica Lag
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Retard des followers par rapport au leader

[ATTENTION]  ALERTE si lag > seuil

JMX Metric :
kafka.server:type=FetcherLagMetrics,name=ConsumerLag,clientId=ReplicaFetcherThread-*


DISASTER RECOVERY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

STRATÉGIE 1 : MIRRORMAKER 2
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Réplication entre clusters Kafka

Cluster Source (Production)     MirrorMaker 2     Cluster Target (DR)
┌─────────────────┐          ──────────────->    ┌─────────────────┐
│ orders (topic)  │                              │ source.orders   │
│ Broker 1,2,3    │                              │ Broker 4,5,6    │
└─────────────────┘                              └─────────────────┘

-> Réplication asynchrone
-> Latence : quelques secondes
-> Préfixe automatique (source.*)

STRATÉGIE 2 : STRETCHED CLUSTER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Un seul cluster sur DEUX datacenters

Datacenter 1                    Datacenter 2
Broker 1, 2                     Broker 3, 4, 5

-> Réplication synchrone
-> Latence plus élevée (inter-DC)
-> Nécessite latence faible entre DCs (<10ms)


BONNES PRATIQUES RÉPLICATION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] replication.factor >= 3 en production
[OK] min.insync.replicas = replication.factor - 1
[OK] acks=all pour données critiques
[OK] unclean.leader.election.enable=false
[OK] Activer rack awareness si multi-rack/DC
[OK] Monitorer under-replicated partitions
[OK] Monitorer offline partitions
[OK] Tester régulièrement les failovers
[OK] Utiliser follower fetching si multi-région
[OK] Preferred leader election activée

[X] Ne PAS utiliser replication.factor=1 en production
[X] Ne PAS ignorer les alertes under-replicated
[X] Ne PAS mettre min.insync.replicas=1 pour données critiques
[X] Ne PAS activer unclean election sans comprendre les risques
"""


# [OK] PARTIE 9 : KAFKA STREAMS (Traitement en temps réel)

"""
┌────────────────────────────────────────────────────────────────────────┐
│              KAFKA STREAMS : STREAM PROCESSING LIBRARY                 │
└────────────────────────────────────────────────────────────────────────┘

QU'EST-CE QUE KAFKA STREAMS ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Bibliothèque Java/Scala pour traiter des flux de données en temps réel

CARACTÉRISTIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Bibliothèque (pas un framework)
[OK] S'intègre dans votre application
[OK] Pas de cluster séparé (contrairement à Spark/Flink)
[OK] Scalabilité horizontale (via partitions)
[OK] Fault-tolerant (state stores répliqués)
[OK] Exactly-once semantics
[OK] Stateful et Stateless operations


COMPARAISON AVEC AUTRES SOLUTIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────────┬──────────────┬──────────────┬──────────────┐
│  Caractéristique │ Kafka Streams│ Spark Stream │ Apache Flink │
├──────────────────┼──────────────┼──────────────┼──────────────┤
│ Type             │ Bibliothèque │ Framework    │ Framework    │
│ Déploiement      │ Application  │ Cluster      │ Cluster      │
│ Langage          │ Java/Scala   │ Multi        │ Java/Scala   │
│ Latence          │ ms           │ secondes     │ ms           │
│ Scalabilité      │ Partitions   │ Executors    │ Task slots   │
│ State            │ RocksDB      │ RDD/DF       │ State backend│
│ Exactly-once     │ Oui [OK]       │ Oui [OK]       │ Oui [OK]       │
│ Complexité       │ Faible       │ Moyenne      │ Élevée       │
│ Kafka natif      │ Oui [OK]       │ Non          │ Non          │
└──────────────────┴──────────────┴──────────────┴──────────────┘


CONCEPTS FONDAMENTAUX
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. STREAM (KStream)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Flux INFINI d'événements

Topic "user-clicks"
┌────────────────────────────────────────────────────────────────┐
│ [Click1] [Click2] [Click3] [Click4] ...                       │
└────────────────────────────────────────────────────────────────┘

Propriétés :
-> INSERT-ONLY (append-only)
-> Chaque enregistrement = nouvel événement
-> Pas de mise à jour en place


2. TABLE (KTable)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Vue MATÉRIALISÉE d'un stream (changelog)

Topic "users" (log compacted)
┌────────────────────────────────────────────────────────────────┐
│ Key=u1, Value={name:"Jean", age:30}    Offset 0               │
│ Key=u2, Value={name:"Marie", age:25}   Offset 1               │
│ Key=u1, Value={name:"Jean", age:31}    Offset 2  <- Update     │
└────────────────────────────────────────────────────────────────┘

État de la KTable à offset 2 :
{
  u1: {name:"Jean", age:31},   <- Dernière valeur pour u1
  u2: {name:"Marie", age:25}
}

Propriétés :
-> UPSERT (insert + update)
-> Chaque clé = dernière valeur
-> Changelog sémantique


3. GLOBALKTABLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Table répliquée sur TOUS les instances

Différence avec KTable :
-> KTable : partitionnée (chaque instance a une partie)
-> GlobalKTable : complète sur chaque instance

Use case :
-> Tables de référence (small datasets)
-> Lookups sans repartitioning


EXEMPLE SIMPLE : WORD COUNT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

// Maven dependency
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>3.6.1</version>
</dependency>

// WordCountApp.java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

import java.util.Arrays;
import java.util.Properties;

public class WordCountApp {
    public static void main(String[] args) {
        
        // 1. CONFIGURATION
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, 
                  Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, 
                  Serdes.String().getClass());
        
        // 2. TOPOLOGIE (DAG de traitement)
        StreamsBuilder builder = new StreamsBuilder();
        
        // Source : topic "text-input"
        KStream<String, String> textLines = builder.stream("text-input");
        
        // Traitement :
        KTable<String, Long> wordCounts = textLines
            // 1. Découper chaque ligne en mots
            .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
            
            // 2. Grouper par mot (repartitionne)
            .groupBy((key, word) -> word)
            
            // 3. Compter les occurrences
            .count();
        
        // Sink : topic "word-count-output"
        wordCounts.toStream().to("word-count-output", 
            Produced.with(Serdes.String(), Serdes.Long()));
        
        // 3. DÉMARRER L'APPLICATION
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        
        // Shutdown graceful
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
        
        streams.start();
        
        System.out.println("Word Count app started [OK]");
    }
}


FLUX DE DONNÉES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Input (topic "text-input") :
"Hello world"
"Hello Kafka Streams"
"Kafka is awesome"

Processing :
1. flatMapValues:
   ["hello", "world"]
   ["hello", "kafka", "streams"]
   ["kafka", "is", "awesome"]

2. groupBy:
   Groupe par mot (repartitioning)

3. count:
   hello  -> 2
   world  -> 1
   kafka  -> 2
   streams-> 1
   is     -> 1
   awesome-> 1

Output (topic "word-count-output") :
hello   2
world   1
kafka   2
streams 1
is      1
awesome 1


OPÉRATIONS STATELESS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. FILTER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Filtrer les enregistrements

KStream<String, Integer> filtered = stream
    .filter((key, value) -> value > 100);

2. MAP / MAPVALUES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Transformer les enregistrements

// map : peut changer clé ET valeur (repartitionne)
KStream<String, String> mapped = stream
    .map((key, value) -> KeyValue.pair(key.toUpperCase(), value * 2));

// mapValues : change seulement valeur (pas de repartitioning)
KStream<String, Integer> mappedValues = stream
    .mapValues(value -> value * 2);

3. FLATMAP / FLATMAPVALUES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Un enregistrement -> 0, 1 ou plusieurs enregistrements

KStream<String, String> flattened = stream
    .flatMapValues(value -> Arrays.asList(value.split(",")));

4. BRANCH
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Diviser un stream selon des prédicats

Map<String, KStream<String, Order>> branches = stream
    .split(Named.as("order-"))
    .branch((key, order) -> order.getAmount() > 1000, Branched.as("high-value"))
    .branch((key, order) -> order.getAmount() > 100, Branched.as("medium-value"))
    .defaultBranch(Branched.as("low-value"));

5. MERGE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Fusionner plusieurs streams

KStream<String, String> merged = stream1.merge(stream2).merge(stream3);


OPÉRATIONS STATEFUL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. AGGREGATIONS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) count()
KTable<String, Long> counts = stream
    .groupByKey()
    .count();

b) reduce()
KTable<String, Integer> sums = stream
    .groupByKey()
    .reduce((aggValue, newValue) -> aggValue + newValue);

c) aggregate()
KTable<String, OrderSummary> aggregated = stream
    .groupByKey()
    .aggregate(
        OrderSummary::new,  // Initializer
        (key, order, summary) -> {  // Aggregator
            summary.addOrder(order);
            return summary;
        },
        Materialized.with(Serdes.String(), orderSummarySerde)
    );


2. WINDOWING (Fenêtrage)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Tumbling Window (fenêtres fixes non-chevauchantes)
┌─────────────────────────────────────────────────────────────────┐
│ Window 1        Window 2        Window 3        Window 4       │
│ [00:00-00:05]   [00:05-00:10]   [00:10-00:15]   [00:15-00:20]  │
│                                                                 │
│ Events:         Events:         Events:         Events:        │
│ E1, E2, E3      E4, E5          E6              E7, E8, E9     │
└─────────────────────────────────────────────────────────────────┘

KTable<Windowed<String>, Long> windowedCounts = stream
    .groupByKey()
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
    .count();


b) Hopping Window (fenêtres glissantes chevauchantes)
┌─────────────────────────────────────────────────────────────────┐
│ Window 1:       [00:00-00:10]                                   │
│ Window 2:           [00:05-00:15]                               │
│ Window 3:               [00:10-00:20]                           │
│                                                                 │
│ -> Fenêtre de 10 min, avance de 5 min                           │
└─────────────────────────────────────────────────────────────────┘

KTable<Windowed<String>, Long> hoppingCounts = stream
    .groupByKey()
    .windowedBy(TimeWindows
        .ofSizeWithNoGrace(Duration.ofMinutes(10))
        .advanceBy(Duration.ofMinutes(5)))
    .count();


c) Session Window (basée sur l'inactivité)
┌─────────────────────────────────────────────────────────────────┐
│ User A:                                                         │
│ Event1  Event2   [GAP 20min]  Event3  Event4                   │
│ └─ Session 1 ──┘               └─ Session 2 ──┘                │
│                                                                 │
│ -> Nouvelle session si gap > inactivityGap                      │
└─────────────────────────────────────────────────────────────────┘

KTable<Windowed<String>, Long> sessionCounts = stream
    .groupByKey()
    .windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(20)))
    .count();


d) Sliding Window
-> Fenêtre continue qui glisse avec chaque événement
-> Utilisé pour operations time-sensitive


3. JOINS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Stream-Stream Join
KStream<String, Order> orders = ...;
KStream<String, Payment> payments = ...;

KStream<String, OrderPayment> joined = orders.join(
    payments,
    (order, payment) -> new OrderPayment(order, payment),
    JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5))
);

-> Join events dans une fenêtre de temps


b) Stream-Table Join (Enrichment)
KStream<String, Order> orders = ...;
KTable<String, Customer> customers = ...;

KStream<String, EnrichedOrder> enriched = orders.join(
    customers,
    (order, customer) -> new EnrichedOrder(order, customer)
);

-> Enrichir un stream avec données d'une table


c) Table-Table Join
KTable<String, User> users = ...;
KTable<String, Address> addresses = ...;

KTable<String, UserWithAddress> joined = users.join(
    addresses,
    (user, address) -> new UserWithAddress(user, address)
);

-> Join co-partitionné de deux tables


STATE STORES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Stockage local (sur disque) de l'état de l'application

TYPES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. KeyValue Store (RocksDB par défaut)
2. Window Store
3. Session Store

CARACTÉRISTIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Persistant (sur disque)
[OK] Fault-tolerant (changelog topic)
[OK] Partitionné (suit le partitioning du input stream)
[OK] Local (pas de RPC)


CHANGELOG TOPIC :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Pour chaque state store, Kafka Streams crée un changelog topic :

Application : "word-count-app"
State store : "counts-store"
Changelog topic : "word-count-app-counts-store-changelog"

┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  STATE STORE (RocksDB local)        CHANGELOG TOPIC           │
│  ┌──────────────────┐               ┌──────────────────┐     │
│  │ hello -> 5        │  ──write──->   │ hello -> 5        │     │
│  │ world -> 3        │               │ world -> 3        │     │
│  │ kafka -> 7        │               │ kafka -> 7        │     │
│  └──────────────────┘               └──────────────────┘     │
│         │                                     │              │
│         │ [IMPACT] CRASH                           │              │
│         [BLACK_DOWN-POINTING_TRIANGLE]                                     │              │
│  ┌──────────────────┐                        │              │
│  │ (vide)           │  [BLACK_LEFT-POINTING_POINTER]──restore───────────┘              │
│  │                  │  replay changelog                     │
│  └──────────────────┘                                       │
└────────────────────────────────────────────────────────────────┘

-> State restauré automatiquement après crash


INTERACTIVE QUERIES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Interroger l'état d'une application Kafka Streams (sans Kafka) :

// Dans l'application Kafka Streams
KTable<String, Long> wordCounts = ...;

// Matérialiser le state store avec un nom
KTable<String, Long> materialized = wordCounts
    .groupBy(...)
    .count(Materialized.as("word-counts-store"));

// API REST pour exposer l'état
@GetMapping("/count/{word}")
public Long getWordCount(@PathVariable String word) {
    ReadOnlyKeyValueStore<String, Long> store = 
        streams.store(
            StoreQueryParameters.fromNameAndType(
                "word-counts-store",
                QueryableStoreTypes.keyValueStore()
            )
        );
    
    return store.get(word);
}

-> Requête O(1) sur le state store local (très rapide)


EXACTLY-ONCE SEMANTICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
          StreamsConfig.EXACTLY_ONCE_V2);

GARANTIES :
[OK] Chaque enregistrement traité EXACTEMENT une fois
[OK] Pas de duplicata dans les aggregations
[OK] State stores cohérents

COMMENT ÇA MARCHE :
-> Transactions Kafka (idempotent producers)
-> Consumer offsets committés transactionnellement
-> Changelog topics transactionnels


EXEMPLE AVANCÉ : E-COMMERCE ANALYTICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Calculer le chiffre d'affaires par catégorie en temps réel :

// Topologie
StreamsBuilder builder = new StreamsBuilder();

// Input : stream d'orders
KStream<String, Order> orders = builder.stream("orders");

// Reference table : products
GlobalKTable<String, Product> products = 
    builder.globalTable("products");

// 1. Enrichir orders avec product info
KStream<String, EnrichedOrder> enriched = orders.join(
    products,
    (orderId, order) -> order.getProductId(),  // Key extractor
    (order, product) -> new EnrichedOrder(order, product)
);

// 2. Grouper par catégorie
KGroupedStream<String, EnrichedOrder> byCategory = enriched
    .groupBy(
        (key, enrichedOrder) -> enrichedOrder.getProduct().getCategory(),
        Grouped.with(Serdes.String(), enrichedOrderSerde)
    );

// 3. Fenêtre glissante de 1 heure
TimeWindowedKStream<String, EnrichedOrder> windowed = byCategory
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofHours(1)));

// 4. Calculer revenu total par fenêtre
KTable<Windowed<String>, Double> revenueByCategory = windowed
    .aggregate(
        () -> 0.0,  // Initializer
        (category, order, revenue) -> revenue + order.getAmount(),
        Materialized.with(Serdes.String(), Serdes.Double())
    );

// 5. Output
revenueByCategory
    .toStream()
    .map((windowedKey, revenue) -> KeyValue.pair(
        windowedKey.key() + "@" + windowedKey.window().start(),
        revenue
    ))
    .to("category-revenue-1h");


TESTING KAFKA STREAMS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TopologyTestDriver pour tester sans Kafka :

@Test
public void testWordCount() {
    // Setup
    Properties props = new Properties();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, "test");
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:1234");
    
    StreamsBuilder builder = new StreamsBuilder();
    // ... build topology
    
    TopologyTestDriver testDriver = new TopologyTestDriver(
        builder.build(), props
    );
    
    // Input topic
    TestInputTopic<String, String> inputTopic = testDriver
        .createInputTopic("input", 
            Serdes.String().serializer(), 
            Serdes.String().serializer());
    
    // Output topic
    TestOutputTopic<String, Long> outputTopic = testDriver
        .createOutputTopic("output",
            Serdes.String().deserializer(),
            Serdes.Long().deserializer());
    
    // Test
    inputTopic.pipeInput("key1", "hello world");
    inputTopic.pipeInput("key2", "hello kafka");
    
    // Assertions
    List<KeyValue<String, Long>> results = outputTopic.readKeyValuesToList();
    assertThat(results).contains(
        KeyValue.pair("hello", 2L),
        KeyValue.pair("world", 1L),
        KeyValue.pair("kafka", 1L)
    );
    
    testDriver.close();
}


MONITORING KAFKA STREAMS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MÉTRIQUES IMPORTANTES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. kafka.streams:type=stream-metrics,client-id=(.*) 
   -> state : état de l'application (RUNNING, REBALANCING, ERROR)

2. kafka.streams:type=stream-task-metrics
   -> process-latency-avg : latence moyenne de traitement
   -> commit-latency-avg : latence de commit

3. kafka.streams:type=stream-processor-node-metrics
   -> process-rate : taux de traitement (records/sec)

4. kafka.streams:type=stream-state-metrics
   -> restore-remaining : records restants à restaurer


BONNES PRATIQUES KAFKA STREAMS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Utiliser exactly-once semantics en production
[OK] Partitionner correctement les input topics
[OK] Utiliser mapValues() plutôt que map() (évite repartitioning)
[OK] Co-partitionner les streams avant join
[OK] Utiliser GlobalKTable pour petites reference tables
[OK] Configurer num.stream.threads selon charge
[OK] Monitorer les métriques
[OK] Tester avec TopologyTestDriver
[OK] Implémenter uncaught exception handler
[OK] Graceful shutdown

[X] Ne PAS utiliser operations bloquantes
[X] Ne PAS faire d'I/O externe dans les transformations
[X] Ne PAS oublier de close() les Serdes custom
[X] Ne PAS négliger le tuning des state stores
"""


# [OK] PARTIE 10 : KAFKA CONNECT (Intégration de données)

"""
┌────────────────────────────────────────────────────────────────────────┐
│              KAFKA CONNECT : DATA INTEGRATION FRAMEWORK                │
└────────────────────────────────────────────────────────────────────────┘

QU'EST-CE QUE KAFKA CONNECT ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Framework distribué pour connecter Kafka avec systèmes externes

SANS KAFKA CONNECT :
┌────────────────────────────────────────────────────────────────┐
│  MySQL ─────-> Script Python ─────-> Kafka Producer            │
│  (You write)      (You maintain)      (You monitor)           │
│                                                                │
│  Problèmes :                                                   │
│  -> Code custom à écrire                                        │
│  -> Maintenance continue                                        │
│  -> Pas de scalabilité automatique                             │
│  -> Pas de fault tolerance                                      │
│  -> Pas de monitoring unifié                                    │
└────────────────────────────────────────────────────────────────┘

AVEC KAFKA CONNECT :
┌────────────────────────────────────────────────────────────────┐
│  MySQL ─────-> JDBC Source Connector ─────-> Kafka              │
│              (Configuration JSON)                              │
│                                                                │
│  Avantages :                                                   │
│  [OK] Pas de code (configuration JSON)                          │
│  [OK] Connectors prêts à l'emploi (centaines)                   │
│  [OK] Scalabilité automatique                                    │
│  [OK] Fault tolerance intégrée                                   │
│  [OK] Monitoring unifié                                          │
│  [OK] Transforms (transformations légères)                       │
└────────────────────────────────────────────────────────────────┘


ARCHITECTURE KAFKA CONNECT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────────────────────────────────────────────────────────┐
│                      KAFKA CONNECT CLUSTER                       │
│                                                                  │
│  ┌──────────────────────────────────────────────────────────┐   │
│  │                    WORKER NODE 1                         │   │
│  │  ┌───────────────────┐    ┌───────────────────┐         │   │
│  │  │ Source Connector  │    │  Sink Connector   │         │   │
│  │  │   (Task 1)        │    │   (Task 1)        │         │   │
│  │  └───────────────────┘    └───────────────────┘         │   │
│  └──────────────────────────────────────────────────────────┘   │
│                                                                  │
│  ┌──────────────────────────────────────────────────────────┐   │
│  │                    WORKER NODE 2                         │   │
│  │  ┌───────────────────┐    ┌───────────────────┐         │   │
│  │  │ Source Connector  │    │  Sink Connector   │         │   │
│  │  │   (Task 2)        │    │   (Task 2)        │         │   │
│  │  └───────────────────┘    └───────────────────┘         │   │
│  └──────────────────────────────────────────────────────────┘   │
│                                                                  │
│  Configuration stockée dans Kafka topics internes :             │
│  -> connect-configs                                               │
│  -> connect-offsets                                               │
│  -> connect-status                                                │
└──────────────────────────────────────────────────────────────────┘
         [BLACK_UP-POINTING_TRIANGLE]                                    │
         │                                    [BLACK_DOWN-POINTING_TRIANGLE]
    ┌─────────┐                         ┌─────────┐
    │  MySQL  │                         │   S3    │
    │  Postgres│                         │  HDFS   │
    │  MongoDB │                         │ Elastic │
    └─────────┘                         └─────────┘


MODES DE DÉPLOIEMENT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. STANDALONE MODE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Un seul worker (développement/test)

bin/connect-standalone.sh config/connect-standalone.properties \
                          config/connector1.properties

Caractéristiques :
-> Configuration dans fichiers .properties
-> Pas de fault tolerance
-> Pas de scalabilité
-> Pour DEV/TEST uniquement

2. DISTRIBUTED MODE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Cluster de workers (production)

bin/connect-distributed.sh config/connect-distributed.properties

Caractéristiques :
[OK] Configuration via REST API
[OK] Fault tolerance automatique
[OK] Scalabilité horizontale
[OK] Rebalancing automatique
[OK] RECOMMANDÉ PRODUCTION


TYPES DE CONNECTORS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. SOURCE CONNECTORS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Importer données VERS Kafka

External System ─────-> Source Connector ─────-> Kafka Topic

Exemples populaires :
-> JDBC Source Connector (MySQL, Postgres, Oracle...)
-> Debezium CDC Connectors (Change Data Capture)
-> MongoDB Source Connector
-> Elasticsearch Source Connector
-> Salesforce Source Connector
-> File Source Connector
-> HTTP Source Connector

2. SINK CONNECTORS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Exporter données DEPUIS Kafka

Kafka Topic ─────-> Sink Connector ─────-> External System

Exemples populaires :
-> JDBC Sink Connector (MySQL, Postgres...)
-> Elasticsearch Sink Connector
-> S3 Sink Connector
-> HDFS Sink Connector
-> MongoDB Sink Connector
-> Cassandra Sink Connector
-> Redis Sink Connector


EXEMPLE 1 : JDBC SOURCE CONNECTOR (MySQL -> Kafka)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
{
  "name": "mysql-source-connector",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": "1",
    
    // Connexion MySQL
    "connection.url": "jdbc:mysql://localhost:3306/mydb",
    "connection.user": "kafka_user",
    "connection.password": "secret",
    
    // Mode de capture
    "mode": "incrementing",
    "incrementing.column.name": "id",
    
    // Tables à surveiller
    "table.whitelist": "users,orders",
    
    // Nom du topic (préfixe)
    "topic.prefix": "mysql-",
    
    // Polling interval
    "poll.interval.ms": "5000"
  }
}

RÉSULTAT :
Table MySQL "users"     -> Topic Kafka "mysql-users"
Table MySQL "orders"    -> Topic Kafka "mysql-orders"

MODES DE CAPTURE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) incrementing
-> Colonne auto-increment (id)
-> Capture seulement NOUVEAUX enregistrements
-> Pas de updates/deletes

b) timestamp
-> Colonne timestamp (updated_at)
-> Capture nouveaux + modifiés
-> Pas de deletes

c) timestamp+incrementing
-> Combine les deux
-> Capture nouveaux + modifiés
-> Pas de deletes

d) bulk (snapshot)
-> Capture tout périodiquement
-> Pas de tracking des changements


EXEMPLE 2 : DEBEZIUM CDC (Change Data Capture)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

POURQUOI DEBEZIUM ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
JDBC Source Connector : Polling (inefficace, pas de deletes)
Debezium : Lit les binlogs/WAL (efficace, tous les changements)

DATABASES SUPPORTÉES :
-> MySQL (binlog)
-> PostgreSQL (logical decoding WAL)
-> MongoDB (oplog)
-> SQL Server (CDC)
-> Oracle (LogMiner)
-> Cassandra

CONFIGURATION DEBEZIUM MYSQL :
{
  "name": "debezium-mysql-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    
    // MySQL connection
    "database.hostname": "localhost",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "secret",
    "database.server.id": "184054",
    "database.server.name": "mydb",
    
    // Tables à surveiller
    "table.include.list": "mydb.users,mydb.orders",
    
    // Snapshot
    "snapshot.mode": "initial",
    
    // Format des messages
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false"
  }
}

RÉSULTAT : Topics créés automatiquement
mydb.users    -> Tous les changements de la table users
mydb.orders   -> Tous les changements de la table orders

FORMAT DES MESSAGES DEBEZIUM :
{
  "before": {  // État AVANT (null pour INSERT)
    "id": 123,
    "name": "John",
    "email": "john@example.com"
  },
  "after": {   // État APRÈS (null pour DELETE)
    "id": 123,
    "name": "John Doe",
    "email": "john.doe@example.com"
  },
  "source": {
    "version": "1.9.0",
    "connector": "mysql",
    "name": "mydb",
    "ts_ms": 1638360000000,
    "snapshot": "false",
    "db": "mydb",
    "table": "users",
    "server_id": 184054,
    "gtid": null,
    "file": "mysql-bin.000003",
    "pos": 154,
    "row": 0
  },
  "op": "u",  // Operation: c=create, u=update, d=delete, r=read (snapshot)
  "ts_ms": 1638360000789
}


EXEMPLE 3 : ELASTICSEARCH SINK CONNECTOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
{
  "name": "elasticsearch-sink",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "tasks.max": "1",
    
    // Elasticsearch connection
    "connection.url": "http://localhost:9200",
    "connection.username": "elastic",
    "connection.password": "secret",
    
    // Topics à consommer
    "topics": "users,orders",
    
    // Index Elasticsearch (un par topic)
    "type.name": "_doc",
    "key.ignore": "false",
    
    // Behavior
    "behavior.on.null.values": "delete",
    "behavior.on.malformed.documents": "warn",
    
    // Batching
    "batch.size": "2000",
    "linger.ms": "1000",
    
    // Schema
    "schema.ignore": "true"
  }
}

FLUX :
Kafka Topic "users" ──-> Elasticsearch Index "users"
  {id:1, name:"John"}      {id:1, name:"John"}


EXEMPLE 4 : S3 SINK CONNECTOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
{
  "name": "s3-sink-connector",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "tasks.max": "4",
    
    // S3 settings
    "s3.region": "us-east-1",
    "s3.bucket.name": "my-kafka-backup",
    
    // Topics
    "topics": "orders,payments",
    
    // Format
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "parquet.codec": "snappy",
    
    // Partitioning (Hive style)
    "partition.duration.ms": "3600000",  // 1 hour
    "path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH",
    "locale": "en",
    "timezone": "UTC",
    
    // Flush
    "flush.size": "1000",
    "rotate.interval.ms": "600000",  // 10 minutes
    
    // Schema
    "schema.compatibility": "NONE"
  }
}

STRUCTURE S3 CRÉÉE :
s3://my-kafka-backup/
  orders/
    year=2024/month=01/day=15/hour=10/
      orders+0+0000000000.parquet
      orders+0+0000001000.parquet
    year=2024/month=01/day=15/hour=11/
      orders+0+0000002000.parquet


SINGLE MESSAGE TRANSFORMS (SMT)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION :
Transformations légères appliquées record par record

TRANSFORMATIONS BUILT-IN :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. InsertField - Ajouter des champs
{
  "transforms": "InsertTimestamp",
  "transforms.InsertTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.InsertTimestamp.timestamp.field": "created_at"
}

2. ReplaceField - Renommer/supprimer des champs
{
  "transforms": "RenameField",
  "transforms.RenameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.RenameField.renames": "user_id:userId,created_at:createdAt",
  "transforms.RenameField.exclude": "internal_field"
}

3. MaskField - Masquer des données sensibles
{
  "transforms": "MaskSSN",
  "transforms.MaskSSN.type": "org.apache.kafka.connect.transforms.MaskField$Value",
  "transforms.MaskSSN.fields": "ssn,credit_card"
}

4. Filter - Filtrer des messages
{
  "transforms": "FilterDeleted",
  "transforms.FilterDeleted.type": "io.confluent.connect.transforms.Filter$Value",
  "transforms.FilterDeleted.filter.condition": "$[?(@.status == 'deleted')]",
  "transforms.FilterDeleted.filter.type": "include"
}

5. ExtractField - Extraire un champ
{
  "transforms": "unwrap",
  "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
  "transforms.unwrap.drop.tombstones": "false"
}

6. TimestampRouter - Router par timestamp
{
  "transforms": "TimestampRouter",
  "transforms.TimestampRouter.type": "org.apache.kafka.connect.transforms.TimestampRouter",
  "transforms.TimestampRouter.topic.format": "${topic}-${timestamp}",
  "transforms.TimestampRouter.timestamp.format": "YYYYMMDD"
}


CHAÎNER PLUSIEURS TRANSFORMS :
{
  "transforms": "InsertTimestamp,RenameFields,MaskSensitive",
  
  "transforms.InsertTimestamp.type": "...",
  "transforms.InsertTimestamp.timestamp.field": "ingested_at",
  
  "transforms.RenameFields.type": "...",
  "transforms.RenameFields.renames": "user_id:userId",
  
  "transforms.MaskSensitive.type": "...",
  "transforms.MaskSensitive.fields": "ssn,email"
}


REST API KAFKA CONNECT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ENDPOINTS PRINCIPAUX :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# 1. Lister les connectors
GET http://localhost:8083/connectors

# 2. Créer un connector
POST http://localhost:8083/connectors
Content-Type: application/json

{
  "name": "my-connector",
  "config": {
    ...
  }
}

# 3. Obtenir le status
GET http://localhost:8083/connectors/my-connector/status

# 4. Mettre en pause
PUT http://localhost:8083/connectors/my-connector/pause

# 5. Redémarrer
POST http://localhost:8083/connectors/my-connector/restart

# 6. Supprimer
DELETE http://localhost:8083/connectors/my-connector

# 7. Obtenir la configuration
GET http://localhost:8083/connectors/my-connector/config

# 8. Mettre à jour la configuration
PUT http://localhost:8083/connectors/my-connector/config
Content-Type: application/json

{
  "connector.class": "...",
  "tasks.max": "2",
  ...
}

# 9. Lister les plugins disponibles
GET http://localhost:8083/connector-plugins

# 10. Valider une configuration
PUT http://localhost:8083/connector-plugins/JdbcSourceConnector/config/validate
Content-Type: application/json

{
  "connector.class": "...",
  ...
}


EXEMPLE COMPLET : Création via curl
curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "mysql-source",
    "config": {
      "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
      "tasks.max": "1",
      "connection.url": "jdbc:mysql://localhost:3306/mydb",
      "connection.user": "kafka",
      "connection.password": "secret",
      "mode": "incrementing",
      "incrementing.column.name": "id",
      "topic.prefix": "mysql-"
    }
  }'


MONITORING KAFKA CONNECT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MÉTRIQUES JMX IMPORTANTES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Connector level
kafka.connect:type=connector-metrics,connector="{connector}"
-> status : état du connector (running, paused, failed)

# Task level
kafka.connect:type=task-metrics,connector="{connector}",task="{task}"
-> status : état de la task

kafka.connect:type=source-task-metrics,connector="{connector}",task="{task}"
-> source-record-poll-rate : taux de polling
-> source-record-write-rate : taux d'écriture dans Kafka

kafka.connect:type=sink-task-metrics,connector="{connector}",task="{task}"
-> sink-record-read-rate : taux de lecture depuis Kafka
-> sink-record-send-rate : taux d'écriture vers externe

# Worker level
kafka.connect:type=connect-worker-metrics
-> task-count : nombre de tasks
-> connector-count : nombre de connectors


BONNES PRATIQUES KAFKA CONNECT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Utiliser distributed mode en production
[OK] Monitorer status et métriques
[OK] Configurer error handling (DLQ)
[OK] Utiliser Avro + Schema Registry
[OK] Tester les transforms avant production
[OK] Dimensionner tasks.max correctement
[OK] Utiliser Debezium pour CDC (pas JDBC polling)
[OK] Configurer offsets.flush.interval.ms
[OK] Backup configuration des connectors
[OK] Utiliser SMT pour transformations simples

[X] Ne PAS utiliser standalone en production
[X] Ne PAS ignorer les erreurs (errors.tolerance=all sans DLQ)
[X] Ne PAS sur-dimensionner tasks.max
[X] Ne PAS faire de transformations lourdes (utiliser Kafka Streams)
[X] Ne PAS oublier de monitorer les DLQ topics
"""


# [OK] PARTIE 11 : SCHEMA REGISTRY (Gestion des schémas)

"""
┌────────────────────────────────────────────────────────────────────────┐
│              SCHEMA REGISTRY : SCHEMA MANAGEMENT                       │
└────────────────────────────────────────────────────────────────────────┘

POURQUOI SCHEMA REGISTRY ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME SANS SCHEMA REGISTRY :
┌────────────────────────────────────────────────────────────────┐
│  Producer V1                        Consumer                   │
│  {                                  ??                         │
│    "id": 123,                       Comment parser ?           │
│    "name": "John",                  Quels champs ?             │
│    "email": "..."                   Quels types ?              │
│  }                                                             │
│                                                                │
│  Producer V2 (breaking change)                                │
│  {                                                             │
│    "userId": 123,  <- Renommé                                  │
│    "fullName": "John Doe",  <- Fusionné                        │
│    "contacts": {...}  <- Nouveau                               │
│  }                                                             │
│                                                                │
│  -> Consumer CRASH [X]                                          │
└────────────────────────────────────────────────────────────────┘

AVEC SCHEMA REGISTRY :
┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  ┌───────────┐       ┌──────────────────┐       ┌─────────┐  │
│  │ Producer  │──1──-> │ Schema Registry  │ <-──2──│Consumer │  │
│  │           │       │ (valide schéma)  │       │         │  │
│  └───────────┘       └──────────────────┘       └─────────┘  │
│       │                       [BLACK_UP-POINTING_TRIANGLE]                       [BLACK_UP-POINTING_TRIANGLE]       │
│       │                       │                       │       │
│       └───────────3───────────┴───────────────────────┘       │
│                          Kafka                                │
│                                                                │
│  1. Producer enregistre schéma                                │
│  2. Consumer récupère schéma                                  │
│  3. Messages envoyés avec schema ID                           │
│                                                                │
│  [OK] Compatibilité garantie                                    │
│  [OK] Évolution contrôlée                                        │
│  [OK] Documentation automatique                                 │
└────────────────────────────────────────────────────────────────┘


EXEMPLE AVEC AVRO
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DÉFINITION DU SCHÉMA :
// user.avsc
{
  "type": "record",
  "name": "User",
  "namespace": "com.example",
  "fields": [
    {"name": "id", "type": "int"},
    {"name": "name", "type": "string"},
    {"name": "email", "type": "string"},
    {"name": "age", "type": ["null", "int"], "default": null}
  ]
}


PRODUCER AVEC SCHEMA REGISTRY (Java) :
import io.confluent.kafka.serializers.KafkaAvroSerializer;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", KafkaAvroSerializer.class.getName());
props.put("schema.registry.url", "http://localhost:8081");

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

// Charger le schéma
String schemaString = "{ ... }";  // Schéma Avro
Schema.Parser parser = new Schema.Parser();
Schema schema = parser.parse(schemaString);

// Créer un enregistrement
GenericRecord user = new GenericData.Record(schema);
user.put("id", 123);
user.put("name", "John Doe");
user.put("email", "john@example.com");
user.put("age", 30);

// Envoyer
ProducerRecord<String, GenericRecord> record = 
    new ProducerRecord<>("users", "user-123", user);
producer.send(record);


CONSUMER AVEC SCHEMA REGISTRY (Java) :
import io.confluent.kafka.serializers.KafkaAvroDeserializer;

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "user-consumer");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", KafkaAvroDeserializer.class.getName());
props.put("schema.registry.url", "http://localhost:8081");

KafkaConsumer<String, GenericRecord> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("users"));

while (true) {
    ConsumerRecords<String, GenericRecord> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, GenericRecord> record : records) {
        GenericRecord user = record.value();
        System.out.println("User ID: " + user.get("id"));
        System.out.println("Name: " + user.get("name"));
        System.out.println("Email: " + user.get("email"));
        System.out.println("Age: " + user.get("age"));
    }
}


ÉVOLUTION DU SCHÉMA (Schema Evolution)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MODES DE COMPATIBILITÉ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. BACKWARD (défaut)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Nouveaux consumers peuvent lire anciennes données

Schéma V1 :                    Schéma V2 (compatible) :
{                              {
  "fields": [                    "fields": [
    {"name": "id"},                {"name": "id"},
    {"name": "name"}               {"name": "name"},
  ]                                {"name": "phone", "default": ""}  <- Nouveau avec default
}                                ]
                               }

[OK] Consumer V2 peut lire messages V1 (utilise default pour phone)
[X] Consumer V1 ne peut PAS lire messages V2

USE CASE : Consumer upgrades first

2. FORWARD
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Anciens consumers peuvent lire nouvelles données

USE CASE : Producer upgrades first

3. FULL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> BACKWARD + FORWARD

[OK] Anciens et nouveaux consumers peuvent lire tout
-> Ajout de champs optionnels avec default uniquement

4. NONE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Pas de vérification de compatibilité


RÈGLES D'ÉVOLUTION AVRO
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

BACKWARD COMPATIBLE :
[OK] Ajouter un champ avec valeur par défaut
[OK] Supprimer un champ avec valeur par défaut
[OK] Changer le type d'un champ vers un type plus général (int -> long)

FORWARD COMPATIBLE :
[OK] Ajouter un champ (consumers anciens l'ignorent)
[OK] Supprimer un champ avec valeur par défaut

NOT COMPATIBLE :
[X] Changer le type d'un champ (incompatible)
[X] Renommer un champ (sans alias)
[X] Supprimer un champ sans default


EXEMPLE D'ÉVOLUTION :
# Version 1
{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "int"},
    {"name": "name", "type": "string"}
  ]
}

# Version 2 (BACKWARD compatible)
{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "int"},
    {"name": "name", "type": "string"},
    {"name": "email", "type": "string", "default": ""}  <- Nouveau avec default [OK]
  ]
}

# Version 3 (BACKWARD compatible)
{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "int"},
    {"name": "name", "type": "string"},
    {"name": "email", "type": "string", "default": ""},
    {"name": "phone", "type": ["null", "string"], "default": null}  <- Nullable [OK]
  ]
}


REST API SCHEMA REGISTRY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# 1. Lister tous les subjects
GET http://localhost:8081/subjects

# 2. Lister les versions d'un subject
GET http://localhost:8081/subjects/users-value/versions

# 3. Obtenir un schéma par version
GET http://localhost:8081/subjects/users-value/versions/1

# 4. Obtenir la dernière version
GET http://localhost:8081/subjects/users-value/versions/latest

# 5. Obtenir un schéma par ID
GET http://localhost:8081/schemas/ids/1

# 6. Enregistrer un nouveau schéma
POST http://localhost:8081/subjects/users-value/versions
Content-Type: application/vnd.schemaregistry.v1+json

{
  "schema": "..."
}

# 7. Vérifier la compatibilité
POST http://localhost:8081/compatibility/subjects/users-value/versions/latest
Content-Type: application/vnd.schemaregistry.v1+json

{
  "schema": "..."
}

Réponse :
{
  "is_compatible": true
}


BONNES PRATIQUES SCHEMA REGISTRY :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Utiliser Avro (compact, performant, schema evolution)
[OK] Choisir BACKWARD ou FULL compatibility
[OK] Toujours fournir des defaults pour nouveaux champs
[OK] Tester la compatibilité avant déploiement
[OK] Versionner les schémas dans Git
[OK] Utiliser des types nullable pour champs optionnels
[OK] Documenter les schémas (doc property)
[OK] Déployer Schema Registry en cluster (HA)
[OK] Monitorer la latence Schema Registry
[OK] Utiliser cache côté client (auto activé)

[X] Ne PAS supprimer des champs sans default
[X] Ne PAS changer les types de manière incompatible
[X] Ne PAS utiliser NONE compatibility en production
[X] Ne PAS renommer des champs sans alias
[X] Ne PAS oublier de tester l'évolution
"""


# [OK] PARTIE 12 : TRANSACTIONS ET EXACTLY-ONCE SEMANTICS

"""
┌────────────────────────────────────────────────────────────────────────┐
│           TRANSACTIONS : EXACTLY-ONCE PROCESSING                       │
└────────────────────────────────────────────────────────────────────────┘

SÉMANTIQUES DE DELIVERY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. AT-MOST-ONCE (Au plus une fois)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Message peut être PERDU
-> Message NE SERA JAMAIS dupliqué

Exemple : Produire sans attendre ACK (acks=0)

Use case : Métriques, logs non-critiques

2. AT-LEAST-ONCE (Au moins une fois)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Message NE SERA JAMAIS perdu
-> Message peut être DUPLIQUÉ

Exemple : Producer avec retries + Consumer commit après traitement

Use case : La plupart des cas (acceptable avec idempotence)

3. EXACTLY-ONCE (Exactement une fois)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Message NE SERA JAMAIS perdu
-> Message NE SERA JAMAIS dupliqué

Use case : Transactions financières, agrégations critiques


PROBLÈME SANS EXACTLY-ONCE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SCÉNARIO : Application de transfert d'argent
┌────────────────────────────────────────────────────────────────┐
│  Topic "payments"          Application          Topic "ledger" │
│                                                                │
│  Payment:                  1. Read payment                     │
│  Transfer $100             2. Process                          │
│  from A to B               3. Write to ledger                  │
│                            4. Commit offset                    │
│                                                                │
│  CRASH après étape 3, avant étape 4 :                         │
│                                                                │
│  -> Payment retraité (offset non committed)                     │
│  -> $100 transféré DEUX FOIS [X]                                │
│  -> Ledger incorrect                                            │
└────────────────────────────────────────────────────────────────┘


TRANSACTIONS KAFKA (Depuis Kafka 0.11)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CARACTÉRISTIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Écriture atomique dans PLUSIEURS partitions
[OK] Read-Committed isolation
[OK] Intégration consumer offsets (commit transactionnel)
[OK] Exactly-once end-to-end


PRODUCER TRANSACTIONNEL (Java)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", StringSerializer.class.getName());

// Configuration TRANSACTIONNELLE
props.put("enable.idempotence", "true");  // Requis
props.put("transactional.id", "my-transactional-producer-1");  // ID unique
props.put("acks", "all");  // Requis
props.put("max.in.flight.requests.per.connection", "5");

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

// Initialiser transactions
producer.initTransactions();

try {
    // Démarrer transaction
    producer.beginTransaction();
    
    // Envoyer plusieurs messages
    producer.send(new ProducerRecord<>("topic1", "key1", "value1"));
    producer.send(new ProducerRecord<>("topic2", "key2", "value2"));
    producer.send(new ProducerRecord<>("topic3", "key3", "value3"));
    
    // Committer TOUTES les écritures atomiquement
    producer.commitTransaction();
    
} catch (ProducerFencedException | OutOfSequenceException | AuthorizationException e) {
    // Erreurs fatales : fermer producer
    producer.close();
} catch (KafkaException e) {
    // Autres erreurs : aborter transaction
    producer.abortTransaction();
}


CONSUMER-PRODUCER TRANSACTIONNEL (Exactly-Once)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Pattern : Read -> Process -> Write (atomique)

// CONSUMER
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "localhost:9092");
consumerProps.put("group.id", "my-consumer-group");
consumerProps.put("key.deserializer", StringDeserializer.class.getName());
consumerProps.put("value.deserializer", StringDeserializer.class.getName());
consumerProps.put("enable.auto.commit", "false");  // Important !
consumerProps.put("isolation.level", "read_committed");  // Important !

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(Collections.singletonList("input-topic"));

// PRODUCER
Properties producerProps = new Properties();
producerProps.put("bootstrap.servers", "localhost:9092");
producerProps.put("key.serializer", StringSerializer.class.getName());
producerProps.put("value.serializer", StringSerializer.class.getName());
producerProps.put("enable.idempotence", "true");
producerProps.put("transactional.id", "my-transformer-1");

KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
producer.initTransactions();

// TRAITEMENT
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    if (!records.isEmpty()) {
        try {
            producer.beginTransaction();
            
            // Traiter et produire
            for (ConsumerRecord<String, String> record : records) {
                String processedValue = processMessage(record.value());
                producer.send(new ProducerRecord<>("output-topic", 
                    record.key(), processedValue));
            }
            
            // Commit offsets DANS la transaction
            Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
            for (TopicPartition partition : records.partitions()) {
                List<ConsumerRecord<String, String>> partitionRecords = 
                    records.records(partition);
                long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
                offsets.put(partition, new OffsetAndMetadata(lastOffset + 1));
            }
            
            producer.sendOffsetsToTransaction(offsets, 
                consumer.groupMetadata());
            
            // Committer ATOMIQUEMENT : écritures + offsets
            producer.commitTransaction();
            
        } catch (ProducerFencedException | OutOfSequenceException | 
                 AuthorizationException e) {
            producer.close();
            consumer.close();
            break;
        } catch (KafkaException e) {
            producer.abortTransaction();
        }
    }
}


ISOLATION LEVEL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. read_uncommitted (défaut)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Lit TOUS les messages (committed + uncommitted)
-> Peut voir des messages d'une transaction aborted

Use case : Performance max, pas besoin d'exactly-once

2. read_committed
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Lit SEULEMENT les messages committed
-> Ignore les messages uncommitted ou aborted
-> Peut avoir de la latence (attend commit)

Use case : Exactly-once required


TRANSACTION COORDINATOR
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FONCTIONNEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Producer envoie beginTransaction() au coordinator
2. Producer écrit dans plusieurs partitions
3. Producer envoie commitTransaction() au coordinator
4. Coordinator écrit COMMIT marker dans TOUTES les partitions
5. Consumers avec read_committed voient les messages

┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  PRODUCER (transactional.id=tx1)                               │
│       │                                                        │
│       │ 1. Begin Transaction                                  │
│       [BLACK_DOWN-POINTING_TRIANGLE]                                                        │
│  ┌─────────────────────────────┐                              │
│  │  Transaction Coordinator    │                              │
│  │  (Broker désigné)           │                              │
│  └─────────────────────────────┘                              │
│       │                                                        │
│       │ 2. Writes                                             │
│       [BLACK_DOWN-POINTING_TRIANGLE]                                                        │
│  ┌─────────┐  ┌─────────┐  ┌─────────┐                       │
│  │ Topic A │  │ Topic B │  │ Offsets │                       │
│  │ Part 0  │  │ Part 1  │  │  Topic  │                       │
│  └─────────┘  └─────────┘  └─────────┘                       │
│       │            │            │                             │
│       │ 3. Commit Transaction   │                             │
│       │            │            │                             │
│       [BLACK_DOWN-POINTING_TRIANGLE]            [BLACK_DOWN-POINTING_TRIANGLE]            [BLACK_DOWN-POINTING_TRIANGLE]                             │
│  [COMMIT]      [COMMIT]    [COMMIT]  <- Markers atomiques     │
│                                                                │
│  Consumers read_committed peuvent maintenant lire [OK]          │
└────────────────────────────────────────────────────────────────┘


TRANSACTION LOG TOPIC
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Topic interne : __transaction_state

Contenu :
-> État des transactions en cours
-> Mapping transactional.id -> Producer Epoch
-> Log compacted

Configuration :
transaction.state.log.replication.factor=3
transaction.state.log.num.partitions=50
transaction.state.log.min.isr=2


PRODUCER EPOCHS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME : Producer zombie (ancien producer pas encore mort)

SOLUTION : EPOCHS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Producer 1 (transactional.id=tx1, epoch=0) démarre
2. Producer 1 crash
3. Producer 2 (transactional.id=tx1, epoch=1) démarre
4. Coordinator incrémente epoch -> 1
5. Producer 1 revient (zombie avec epoch=0)
6. Producer 1 essaie d'écrire
7. Broker REJETTE (epoch=0 < epoch=1)

-> ProducerFencedException levée
-> Zombie fencé automatiquement [OK]


IDEMPOTENT PRODUCER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

enable.idempotence=true

GARANTIES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Pas de duplicatas (même si retry)
[OK] Ordre strict (même avec retries)
[OK] Par partition

FONCTIONNEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque message a :
-> Producer ID (PID) - assigné par broker
-> Sequence Number - incrémenté par partition

Broker maintient :
-> Last Sequence Number par (PID, Partition)

Si message reçu avec Sequence Number déjà vu :
-> ACK renvoyé (succès)
-> Message PAS écrit (dédupliqué) [OK]

Configuration automatique avec enable.idempotence=true :
-> acks=all
-> max.in.flight.requests.per.connection <= 5
-> retries > 0


KAFKA STREAMS EXACTLY-ONCE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
          StreamsConfig.EXACTLY_ONCE_V2);

// Ou legacy (Kafka < 2.5)
StreamsConfig.EXACTLY_ONCE

GARANTIES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Exactly-once end-to-end
[OK] State stores cohérents
[OK] Pas de duplicatas dans aggregations
[OK] Changelog topics transactionnels

EXEMPLE :
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
          StreamsConfig.EXACTLY_ONCE_V2);

StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> textLines = builder.stream("input");

KTable<String, Long> wordCounts = textLines
    .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
    .groupBy((key, word) -> word)
    .count();

wordCounts.toStream().to("output");

KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

-> Chaque mot compté EXACTEMENT une fois [OK]


PERFORMANCES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

IMPACT DES TRANSACTIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Latence :
-> +5-10ms en moyenne (commit markers)
-> Dépend de la taille des transactions

Throughput :
-> Légère baisse (5-10%)
-> Plus de metadata à stocker

Recommandations :
[OK] Batcher les messages dans une transaction (pas 1 msg/transaction)
[OK] Utiliser max 1000 messages par transaction
[OK] Configurer transaction.max.timeout.ms appropriément


BONNES PRATIQUES TRANSACTIONS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Utiliser exactly-once pour données critiques
[OK] transactional.id unique par producer instance
[OK] enable.idempotence=true (toujours)
[OK] isolation.level=read_committed pour consumers
[OK] Batcher les messages dans transactions
[OK] Gérer ProducerFencedException correctement
[OK] Monitorer transaction.max.timeout.ms
[OK] Kafka Streams : EXACTLY_ONCE_V2 (plus performant)

[X] Ne PAS réutiliser transactional.id sur plusieurs instances
[X] Ne PAS faire de transactions trop longues (timeout)
[X] Ne PAS oublier abortTransaction() dans catch
[X] Ne PAS utiliser auto-commit avec transactions
[X] Ne PAS mélanger transactional et non-transactional dans même topic
"""


# [OK] PARTIE 13 : PERFORMANCE, TUNING ET MONITORING

"""
┌────────────────────────────────────────────────────────────────────────┐
│              PERFORMANCE TUNING & MONITORING                           │
└────────────────────────────────────────────────────────────────────────┘

MÉTRIQUES CLÉS À MONITORER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. BROKER METRICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Under Replicated Partitions [ATTENTION]  CRITIQUE
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions

-> Partitions avec répliques non synchronisées
-> ALERTE si > 0
-> Causes : Broker lent, network issues, I/O problems

b) Offline Partitions [ATTENTION]  CRITIQUE
kafka.controller:type=KafkaController,name=OfflinePartitionsCount

-> Partitions sans leader (indisponibles)
-> ALERTE si > 0
-> Intervention immédiate requise

c) Active Controller Count
kafka.controller:type=KafkaController,name=ActiveControllerCount

-> Doit être exactement 1 dans le cluster
-> ALERTE si 0 ou > 1

d) Request Handler Avg Idle Percent
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent

-> % de temps idle des request handlers
-> < 30% = surcharge
-> Augmenter num.network.threads ou num.io.threads

e) Bytes In/Out Per Second
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec
kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec

-> Throughput du broker
-> Monitorer pour capacity planning

f) Messages In Per Second
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec

-> Taux de messages entrants
-> Identifier pics de charge

g) Failed Produce/Fetch Requests
kafka.server:type=BrokerTopicMetrics,name=FailedProduceRequestsPerSec
kafka.server:type=BrokerTopicMetrics,name=FailedFetchRequestsPerSec

-> ALERTE si > 0
-> Investiguer causes (quotas, auth, timeouts)

h) ISR Shrinks/Expands
kafka.server:type=ReplicaManager,name=IsrShrinksPerSec
kafka.server:type=ReplicaManager,name=IsrExpandsPerSec

-> Fréquence de changement ISR
-> Trop fréquent = instabilité


2. PRODUCER METRICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Record Send Rate
kafka.producer:type=producer-metrics,client-id=<id>,attribute=record-send-rate

-> Messages/sec produits
-> Baseline pour monitoring

b) Record Error Rate
kafka.producer:type=producer-metrics,client-id=<id>,attribute=record-error-rate

-> ALERTE si > 0
-> Vérifier logs

c) Record Retry Rate
kafka.producer:type=producer-metrics,client-id=<id>,attribute=record-retry-rate

-> Taux de retries
-> Élevé = problème network/broker

d) Request Latency Avg/Max
kafka.producer:type=producer-metrics,client-id=<id>,attribute=request-latency-avg
kafka.producer:type=producer-metrics,client-id=<id>,attribute=request-latency-max

-> Latence des requests
-> Baseline : 2-10ms (réseau local)

e) Buffer Available Bytes
kafka.producer:type=producer-metrics,client-id=<id>,attribute=buffer-available-bytes

-> Espace buffer disponible
-> 0 = backpressure (ralentit producer)

f) Compression Rate
kafka.producer:type=producer-topic-metrics,client-id=<id>,topic=<topic>,attribute=compression-rate

-> Ratio de compression
-> Vérifier efficacité compression


3. CONSUMER METRICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Consumer Lag [ATTENTION]  IMPORTANT
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=<id>,partition=<partition>,topic=<topic>,attribute=records-lag

-> Retard du consumer (messages non consommés)
-> ALERTE si lag croissant
-> Actions : scaler consumers, optimiser traitement

b) Records Consumed Rate
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=<id>,attribute=records-consumed-rate

-> Messages/sec consommés
-> Comparer avec producer rate

c) Fetch Latency Avg/Max
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=<id>,attribute=fetch-latency-avg

-> Latence des fetch requests
-> Élevée = broker surchargé

d) Commit Latency Avg
kafka.consumer:type=consumer-coordinator-metrics,client-id=<id>,attribute=commit-latency-avg

-> Latence des commits
-> Élevée = coordinator surchargé


TUNING PRODUCER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF : THROUGHPUT MAX
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

batch.size=32768  // 32KB (défaut 16KB)
linger.ms=10      // Attendre 10ms pour remplir batch
compression.type=zstd  // Meilleur compromis
buffer.memory=67108864  // 64MB (défaut 32MB)
max.in.flight.requests.per.connection=5

Résultat :
-> +200% throughput
-> +5-10ms latence

OBJECTIF : LATENCE MIN
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

batch.size=1      // Pas de batching
linger.ms=0       // Envoyer immédiatement
compression.type=none
max.in.flight.requests.per.connection=1

Résultat :
-> Latence ~1-2ms
-> Throughput divisé par 10

OBJECTIF : DURABILITÉ MAX
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

acks=all
enable.idempotence=true
max.in.flight.requests.per.connection=5
retries=Integer.MAX_VALUE
retry.backoff.ms=100

Résultat :
-> Pas de perte de données
-> Latence +10-20ms


TUNING CONSUMER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF : THROUGHPUT MAX
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

fetch.min.bytes=1048576  // 1MB (attendre plus de données)
fetch.max.wait.ms=500    // Max 500ms d'attente
max.poll.records=1000    // 1000 messages par poll
max.partition.fetch.bytes=2097152  // 2MB par partition

Résultat :
-> Moins de fetch requests
-> Meilleur throughput
-> +latence si peu de messages

OBJECTIF : LATENCE MIN
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

fetch.min.bytes=1       // Ne pas attendre
fetch.max.wait.ms=0     // Retourner immédiatement
max.poll.records=100    // Petits batches

Résultat :
-> Latence minimale
-> Plus de fetch requests
-> Moins de throughput


TUNING BROKER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

THREADS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

num.network.threads=8  // Threads réseau (1 par core)
num.io.threads=16      // Threads I/O disque (2x cores)
num.replica.fetchers=4 // Threads réplication

MÉMOIRE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# JVM Heap
export KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"

Règle : 6GB heap, reste pour page cache OS

# Socket buffers
socket.send.buffer.bytes=1048576    // 1MB
socket.receive.buffer.bytes=1048576 // 1MB
socket.request.max.bytes=104857600  // 100MB

LOG RETENTION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

log.retention.hours=168          // 7 jours
log.retention.bytes=-1           // Illimité
log.segment.bytes=1073741824     // 1GB
log.roll.hours=168               // Nouveau segment tous les 7j

FLUSH POLICY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

log.flush.interval.messages=10000
log.flush.interval.ms=1000

[ATTENTION]  Recommandation : Laisser OS gérer flush (meilleures performances)
-> Compter sur réplication pour durabilité

RÉPLICATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

default.replication.factor=3
min.insync.replicas=2
replica.lag.time.max.ms=30000
replica.fetch.max.bytes=1048576


TUNING OS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FILE DESCRIPTORS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# /etc/security/limits.conf
kafka soft nofile 100000
kafka hard nofile 100000

Vérifier :
ulimit -n

VM SETTINGS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# /etc/sysctl.conf
vm.swappiness=1                    # Minimiser swap
vm.dirty_background_ratio=5        # Flush asynchrone à 5%
vm.dirty_ratio=80                  # Flush synchrone à 80%
vm.max_map_count=262144            # Pour grands index

Appliquer :
sudo sysctl -p

DISK I/O SCHEDULER
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Pour SSD
echo noop > /sys/block/sda/queue/scheduler

# Pour HDD
echo deadline > /sys/block/sda/queue/scheduler

FILESYSTEM
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Recommandé : XFS ou EXT4

Mount options :
noatime,nodiratime  # Pas de update access time


MONITORING STACK
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARCHITECTURE RECOMMANDÉE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  Kafka Cluster                                                 │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐                    │
│  │ Broker 1 │  │ Broker 2 │  │ Broker 3 │                    │
│  │ (JMX)    │  │ (JMX)    │  │ (JMX)    │                    │
│  └─────┬────┘  └─────┬────┘  └─────┬────┘                    │
│        │             │             │                          │
│        └─────────────┴─────────────┘                          │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │ JMX Exporter     │                               │
│            │ (Prometheus)     │                               │
│            └─────────┬────────┘                               │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │   Prometheus     │                               │
│            │   (Time Series)  │                               │
│            └─────────┬────────┘                               │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │    Grafana       │                               │
│            │  (Dashboards)    │                               │
│            └──────────────────┘                               │
└────────────────────────────────────────────────────────────────┘


JMX EXPORTER CONFIGURATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# kafka-server-start.sh
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false \
  -javaagent:/opt/jmx_prometheus_javaagent.jar=7071:/opt/kafka-broker.yml"

# kafka-broker.yml (JMX Exporter rules)
lowercaseOutputName: true
rules:
  - pattern: kafka.server<type=(.+), name=(.+)><>Value
    name: kafka_server_$1_$2
  - pattern: kafka.controller<type=(.+), name=(.+)><>Value
    name: kafka_controller_$1_$2


PROMETHEUS CONFIGURATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# prometheus.yml
scrape_configs:
  - job_name: 'kafka'
    static_configs:
      - targets:
        - 'broker1:7071'
        - 'broker2:7071'
        - 'broker3:7071'


GRAFANA DASHBOARDS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Dashboards populaires :
-> Kafka Exporter Overview (ID: 7589)
-> Kafka Overview (ID: 12460)
-> Kafka Cluster (ID: 14014)

Importer depuis https://grafana.com/grafana/dashboards/


ALERTING RULES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# alert.rules.yml
groups:
  - name: kafka
    rules:
      # Under Replicated Partitions
      - alert: KafkaUnderReplicatedPartitions
        expr: kafka_server_replicamanager_underreplicatedpartitions > 0
        for: 5m
        labels:
          severity: critical
        annotations:
          summary: "Kafka has under-replicated partitions"
          description: "Broker {{ $labels.instance }} has {{ $value }} under-replicated partitions"
      
      # Offline Partitions
      - alert: KafkaOfflinePartitions
        expr: kafka_controller_kafkacontroller_offlinepartitionscount > 0
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: "Kafka has offline partitions"
          description: "{{ $value }} partitions are offline"
      
      # Consumer Lag
      - alert: KafkaConsumerLag
        expr: kafka_consumer_lag > 10000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High consumer lag"
          description: "Consumer {{ $labels.group }} has lag of {{ $value }} on topic {{ $labels.topic }}"
      
      # No Active Controller
      - alert: KafkaNoActiveController
        expr: kafka_controller_kafkacontroller_activecontrollercount != 1
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: "No active Kafka controller"
          description: "Active controller count is {{ $value }}"


CAPACITY PLANNING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CALCUL DU NOMBRE DE BROKERS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Facteurs :
1. Throughput requis (MB/s)
2. Retention (jours)
3. Replication factor
4. Capacité disque par broker

Formule :
Brokers = MAX(
  Throughput_requis / Throughput_par_broker,
  (Data_journaliere * Retention * Replication) / Capacité_disque
)

Exemple :
-> Throughput : 500 MB/s
-> Retention : 7 jours
-> Replication : 3
-> Throughput/broker : 50 MB/s
-> Disque/broker : 2 TB

Calcul throughput :
500 / 50 = 10 brokers

Calcul stockage :
Data/jour = 500 MB/s * 86400s = 43.2 TB/jour
Storage total = 43.2 * 7 * 3 = 907.2 TB
Brokers = 907.2 / 2 = 454 brokers

-> Résultat : 454 brokers (limité par stockage)


CALCUL DU NOMBRE DE PARTITIONS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Règle générale :
Partitions = MAX(
  Nombre de consumers,
  Throughput_requis / Throughput_par_partition
)

Limites :
-> Max ~4000 partitions par broker
-> Latence augmente avec nombre de partitions
-> Election leader plus lente

Recommandation :
-> Commencer conservateur (10-50 partitions)
-> Augmenter si nécessaire
-> Impossible de réduire


BONNES PRATIQUES PERFORMANCE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Monitorer Under Replicated Partitions
[OK] Monitorer Consumer Lag
[OK] Utiliser compression (zstd recommandé)
[OK] Batcher les messages (producer + consumer)
[OK] SSD pour meilleure performance
[OK] Dimensionner correctement le heap JVM
[OK] Tuner OS (file descriptors, vm.swappiness)
[OK] Page cache > Heap
[OK] Utiliser Grafana pour dashboards
[OK] Alerter sur métriques critiques
[OK] Capacity planning régulier
[OK] Tester les performances (kafka-producer-perf-test)

[X] Ne PAS ignorer Under Replicated Partitions
[X] Ne PAS sous-dimensionner les brokers
[X] Ne PAS créer trop de partitions d'un coup
[X] Ne PAS négliger le monitoring
[X] Ne PAS utiliser RAID pour disques (JBOD préféré)
[X] Ne PAS mettre trop de heap (6-8GB max)
"""


# [OK] PARTIE 14 : CAS D'USAGE RÉELS

"""
┌────────────────────────────────────────────────────────────────────────┐
│              CAS D'USAGE KAFKA DANS L'INDUSTRIE                        │
└────────────────────────────────────────────────────────────────────────┘

CAS #1 : LINKEDIN - ACTIVITY TRACKING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> 900+ millions d'utilisateurs
-> 7+ trillions de messages par jour
-> Activity streams, metrics, logs

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│  LinkedIn Apps                                                 │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐                    │
│  │  Web     │  │  Mobile  │  │   API    │                    │
│  └─────┬────┘  └─────┬────┘  └─────┬────┘                    │
│        │             │             │                          │
│        └─────────────┴─────────────┘                          │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │  Kafka Cluster   │                               │
│            │  (1000+ brokers) │                               │
│            └─────────┬────────┘                               │
│                      │                                        │
│        ┌─────────────┼─────────────┐                          │
│        │             │             │                          │
│        [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]                          │
│  ┌──────────┐ ┌──────────┐ ┌──────────┐                      │
│  │Analytics │ │  Hadoop  │ │Real-time │                      │
│  │          │ │   ETL    │ │ Metrics  │                      │
│  └──────────┘ └──────────┘ └──────────┘                      │
└────────────────────────────────────────────────────────────────┘

TOPICS PRINCIPAUX :
-> page-views : Navigation utilisateurs
-> impressions : Affichages publicités
-> connections : Nouvelles connexions
-> messages : Messages entre utilisateurs

RÉSULTATS :
[OK] 7+ trillions messages/jour
[OK] Latence < 10ms
[OK] 99.999% availability


CAS #2 : UBER - REAL-TIME PRICING & ETA
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> Surge pricing dynamique
-> ETA en temps réel
-> Matching riders-drivers

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  Rider App ─────┐                                             │
│                 │                                             │
│  Driver App ────┼──-> Kafka Topics                            │
│                 │    ┌─────────────────┐                     │
│  GPS Sensors ───┘    │ driver-locations│                     │
│                      │ ride-requests   │                     │
│                      │ surge-events    │                     │
│                      └────────┬────────┘                     │
│                               │                              │
│                ┌──────────────┼──────────────┐               │
│                │              │              │               │
│                [BLACK_DOWN-POINTING_TRIANGLE]              [BLACK_DOWN-POINTING_TRIANGLE]              [BLACK_DOWN-POINTING_TRIANGLE]               │
│         ┌──────────┐   ┌──────────┐   ┌──────────┐          │
│         │ Kafka    │   │ Flink    │   │ Spark    │          │
│         │ Streams  │   │Streaming │   │Streaming │          │
│         │          │   │          │   │          │          │
│         │ Matching │   │ Pricing  │   │Analytics │          │
│         └──────────┘   └──────────┘   └──────────┘          │
│                │              │              │               │
│                └──────────────┼──────────────┘               │
│                               [BLACK_DOWN-POINTING_TRIANGLE]                              │
│                        ┌──────────────┐                      │
│                        │ Redis Cache  │                      │
│                        │ (Driver ETA) │                      │
│                        └──────────────┘                      │
└────────────────────────────────────────────────────────────────┘

TOPICS :
-> driver-locations (GPS updates, 1Hz)
-> ride-requests (Demandes courses)
-> surge-events (Calculs tarification)
-> trips (Courses en cours/terminées)

KAFKA STREAMS PROCESSING :
-> Matching algorithm (rider <-> driver)
-> ETA calculation (temps trajet)
-> Surge pricing (offre/demande)

RÉSULTATS :
[OK] Matching en < 1 seconde
[OK] ETA précis (erreur < 5%)
[OK] Surge pricing réactif (secondes)


CAS #3 : NETFLIX - REAL-TIME RECOMMENDATIONS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> 200+ millions abonnés
-> Recommandations personnalisées
-> A/B testing en temps réel

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│  User Events                                                   │
│  ┌────────────────────────────────────────────────┐            │
│  │ play, pause, rewind, rate, search, browse...  │            │
│  └─────────────────────┬──────────────────────────┘            │
│                        │                                       │
│                        [BLACK_DOWN-POINTING_TRIANGLE]                                       │
│              ┌──────────────────┐                              │
│              │  Kafka Topics    │                              │
│              │  - user-events   │                              │
│              │  - viewing-data  │                              │
│              └────────┬─────────┘                              │
│                       │                                        │
│         ┌─────────────┼─────────────┐                          │
│         │             │             │                          │
│         [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]                          │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐                     │
│  │  Kafka   │  │   Flink  │  │  Spark   │                     │
│  │ Streams  │  │Streaming │  │ ML Jobs  │                     │
│  │          │  │          │  │          │                     │
│  │Real-time │  │Trending  │  │  Model   │                     │
│  │Analytics │  │Detection │  │ Training │                     │
│  └─────┬────┘  └─────┬────┘  └─────┬────┘                     │
│        │             │             │                           │
│        └─────────────┴─────────────┘                           │
│                      │                                         │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                         │
│            ┌──────────────────┐                                │
│            │ Recommendation   │                                │
│            │    Service       │                                │
│            │  (Cassandra)     │                                │
│            └──────────────────┘                                │
└────────────────────────────────────────────────────────────────┘

PROCESSING :
-> Event aggregation (vues par catégorie)
-> Trending detection (contenus populaires)
-> User profiling (préférences)
-> A/B test scoring (expériences)

RÉSULTATS :
[OK] 80% du contenu visionné via recommandations
[OK] Latence recommandations < 100ms
[OK] A/B testing instantané


CAS #4 : AIRBNB - PAYMENT PROCESSING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> Transactions financières (réservations, paiements)
-> Exactly-once semantics requis
-> Audit trail complet

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│  Booking Service                                               │
│       │                                                        │
│       │ 1. Create booking                                     │
│       [BLACK_DOWN-POINTING_TRIANGLE]                                                        │
│  ┌──────────────────┐                                         │
│  │ Kafka Topics     │                                         │
│  │ (Transactions)   │                                         │
│  │                  │                                         │
│  │ - bookings       │                                         │
│  │ - payments       │                                         │
│  │ - refunds        │                                         │
│  └────────┬─────────┘                                         │
│           │                                                   │
│           │ 2. Transactional processing                      │
│           [BLACK_DOWN-POINTING_TRIANGLE]                                                   │
│  ┌──────────────────┐                                         │
│  │  Kafka Streams   │                                         │
│  │  (Exactly-Once)  │                                         │
│  │                  │                                         │
│  │  Payment Logic   │                                         │
│  └────────┬─────────┘                                         │
│           │                                                   │
│           ├──-> Payment Gateway (Stripe/Adyen)                 │
│           ├──-> Accounting System                              │
│           ├──-> Fraud Detection                                │
│           └──-> Audit Log                                      │
└────────────────────────────────────────────────────────────────┘

EXACTLY-ONCE PROCESSING :
-> enable.idempotence=true
-> transactional.id unique par instance
-> isolation.level=read_committed

TOPICS :
-> bookings : Nouvelles réservations
-> payments : Paiements traités
-> refunds : Remboursements
-> audit-log : Log complet (compliance)

RÉSULTATS :
[OK] 0 duplicatas paiements
[OK] Audit trail complet
[OK] Compliance financière


CAS #5 : WALMART - INVENTORY MANAGEMENT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> 10,000+ magasins
-> Millions de produits
-> Stock en temps réel

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│  Point of Sale (10,000+ stores)                                │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐                    │
│  │ Store 1  │  │ Store 2  │  │Store N..││                    │
│  │ POS      │  │ POS      │  │ POS      │                    │
│  └─────┬────┘  └─────┬────┘  └─────┬────┘                    │
│        │             │             │                          │
│        └─────────────┴─────────────┘                          │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │  Kafka Connect   │                               │
│            │  (CDC Debezium)  │                               │
│            └─────────┬────────┘                               │
│                      │                                        │
│                      [BLACK_DOWN-POINTING_TRIANGLE]                                        │
│            ┌──────────────────┐                               │
│            │  Kafka Topics    │                               │
│            │  - sales         │                               │
│            │  - inventory     │                               │
│            └─────────┬────────┘                               │
│                      │                                        │
│        ┌─────────────┼─────────────┐                          │
│        │             │             │                          │
│        [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]                          │
│  ┌──────────┐ ┌──────────┐ ┌──────────┐                      │
│  │ Kafka    │ │ Elastic  │ │  Data    │                      │
│  │ Streams  │ │  Search  │ │  Lake    │                      │
│  │          │ │          │ │ (S3/HDFS)│                      │
│  │Aggregate │ │Real-time │ │Analytics │                      │
│  │ Stock    │ │  Search  │ │          │                      │
│  └─────┬────┘ └──────────┘ └──────────┘                      │
│        │                                                      │
│        [BLACK_DOWN-POINTING_TRIANGLE]                                                      │
│  ┌──────────────────┐                                         │
│  │  Cassandra DB    │                                         │
│  │  (Stock levels)  │                                         │
│  └──────────────────┘                                         │
└────────────────────────────────────────────────────────────────┘

PROCESSING :
-> CDC depuis POS databases (Debezium)
-> Aggregate stock par produit/store
-> Alertes rupture de stock
-> Prédictions réapprovisionnement

RÉSULTATS :
[OK] Stock temps réel (< 5 secondes)
[OK] Réduction ruptures de stock (-30%)
[OK] Optimisation approvisionnement


CAS #6 : TWITTER - TWEET PROCESSING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTEXTE :
-> 500+ millions tweets par jour
-> Timeline generation
-> Trending topics

ARCHITECTURE :
┌────────────────────────────────────────────────────────────────┐
│  Twitter API                                                   │
│       │                                                        │
│       │ Post tweet                                            │
│       [BLACK_DOWN-POINTING_TRIANGLE]                                                        │
│  ┌──────────────────┐                                         │
│  │  Kafka Topics    │                                         │
│  │  - tweets        │                                         │
│  │  - retweets      │                                         │
│  │  - likes         │                                         │
│  └────────┬─────────┘                                         │
│           │                                                   │
│     ┌─────┼─────┬─────────────┐                               │
│     │     │     │             │                               │
│     [BLACK_DOWN-POINTING_TRIANGLE]     [BLACK_DOWN-POINTING_TRIANGLE]     [BLACK_DOWN-POINTING_TRIANGLE]             [BLACK_DOWN-POINTING_TRIANGLE]                               │
│  ┌────┐┌────┐┌──────┐  ┌──────────┐                          │
│  │Fan-││Trend││Search│  │Analytics │                          │
│  │out ││Calc ││Index │  │  (Hadoop)│                          │
│  └──┬─┘└────┘└──────┘  └──────────┘                          │
│     │                                                         │
│     └──-> Timeline Service                                     │
└────────────────────────────────────────────────────────────────┘

FANOUT PROCESSING :
Pour chaque tweet :
-> Trouver tous les followers
-> Publier dans timeline de chaque follower
-> Utilise Kafka pour scalabilité

RÉSULTATS :
[OK] 500M+ tweets/jour traités
[OK] Timeline < 200ms
[OK] Trending topics temps réel


CAS #7 : E-COMMERCE - ORDER PROCESSING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARCHITECTURE MICROSERVICES :
┌────────────────────────────────────────────────────────────────┐
│                                                                │
│  Web/Mobile App                                                │
│       │                                                        │
│       │ Place order                                           │
│       [BLACK_DOWN-POINTING_TRIANGLE]                                                        │
│  ┌──────────────────┐                                         │
│  │ Order Service    │───-> orders topic                        │
│  └──────────────────┘                                         │
│                            │                                  │
│         ┌──────────────────┼──────────────────┬───────────┐   │
│         │                  │                  │           │   │
│         [BLACK_DOWN-POINTING_TRIANGLE]                  [BLACK_DOWN-POINTING_TRIANGLE]                  [BLACK_DOWN-POINTING_TRIANGLE]           [BLACK_DOWN-POINTING_TRIANGLE]   │
│  ┌────────────┐    ┌────────────┐    ┌────────────┐┌────────┐│
│  │ Payment    │    │ Inventory  │    │  Shipping  ││Email   ││
│  │ Service    │    │  Service   │    │  Service   ││Service ││
│  └──────┬─────┘    └──────┬─────┘    └──────┬─────┘└───┬────┘│
│         │                  │                  │          │    │
│         └──────────────────┴──────────────────┴──────────┘    │
│                            │                                  │
│                            [BLACK_DOWN-POINTING_TRIANGLE]                                  │
│                    ┌────────────────┐                         │
│                    │ Order Status   │                         │
│                    │    (Updated)   │                         │
│                    └────────────────┘                         │
└────────────────────────────────────────────────────────────────┘

SAGA PATTERN :
1. Order created -> orders topic
2. Payment processed -> payments topic
3. Inventory reserved -> inventory topic
4. Shipping created -> shipping topic
5. Email sent -> notifications topic

En cas d'erreur (ex: payment failed) :
-> Compensation events (rollback)
-> Cancel inventory, shipping, etc.

RÉSULTATS :
[OK] Découplage microservices
[OK] Résilience (retry automatique)
[OK] Audit trail complet
[OK] Scalabilité indépendante
"""


# [OK] PARTIE 15 : EXEMPLES PRATIQUES COMPLETS

"""
┌────────────────────────────────────────────────────────────────────────┐
│              EXEMPLES PRATIQUES COMPLETS                               │
└────────────────────────────────────────────────────────────────────────┘

EXEMPLE #1 : SYSTÈME DE LOGS CENTRALISÉ
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF :
Agréger logs de 100+ microservices vers Elasticsearch

ARCHITECTURE :
Microservices -> Filebeat -> Kafka -> Logstash -> Elasticsearch -> Kibana

CONFIGURATION FILEBEAT :
# filebeat.yml
filebeat.inputs:
  - type: log
    enabled: true
    paths:
      - /var/log/myapp/*.log
    fields:
      service: myapp
      environment: production

output.kafka:
  hosts: ["kafka1:9092", "kafka2:9092", "kafka3:9092"]
  topic: 'logs-%{[fields.service]}'
  partition.round_robin:
    reachable_only: false
  required_acks: 1
  compression: gzip
  max_message_bytes: 1000000

CONFIGURATION LOGSTASH :
# logstash.conf
input {
  kafka {
    bootstrap_servers => "kafka1:9092,kafka2:9092,kafka3:9092"
    topics_pattern => "logs-.*"
    group_id => "logstash-consumers"
    consumer_threads => 4
    codec => json
  }
}

filter {
  # Parse JSON logs
  json {
    source => "message"
  }
  
  # Add timestamp
  date {
    match => [ "timestamp", "ISO8601" ]
    target => "@timestamp"
  }
  
  # Enrichissement
  mutate {
    add_field => { "processed_at" => "%{@timestamp}" }
  }
}

output {
  elasticsearch {
    hosts => ["es1:9200", "es2:9200", "es3:9200"]
    index => "logs-%{[fields][service]}-%{+YYYY.MM.dd}"
  }
}

KAFKA TOPICS CRÉATION :
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create \
  --topic logs-api \
  --partitions 10 \
  --replication-factor 3 \
  --config retention.ms=604800000  # 7 jours

RÉSULTATS :
[OK] 1M+ logs/seconde
[OK] Latence end-to-end < 2 secondes
[OK] Aucune perte de logs (réplication)


EXEMPLE #2 : REAL-TIME ANALYTICS DASHBOARD
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF :
Dashboard temps réel des ventes e-commerce

ARCHITECTURE :
Web App -> Kafka -> Kafka Streams -> Websocket -> Dashboard

PRODUCER (Web App) :
// SalesProducer.java
import org.apache.kafka.clients.producer.*;
import com.fasterxml.jackson.databind.ObjectMapper;

public class SalesProducer {
    
    private final KafkaProducer<String, String> producer;
    private final ObjectMapper mapper = new ObjectMapper();
    
    public SalesProducer() {
        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");
        props.put("acks", "1");
        props.put("compression.type", "snappy");
        
        this.producer = new KafkaProducer<>(props);
    }
    
    public void sendSale(Sale sale) throws Exception {
        String saleJson = mapper.writeValueAsString(sale);
        ProducerRecord<String, String> record = 
            new ProducerRecord<>("sales", sale.getProductId(), saleJson);
        
        producer.send(record, (metadata, exception) -> {
            if (exception != null) {
                exception.printStackTrace();
            }
        });
    }
}

// Sale.java
public class Sale {
    private String orderId;
    private String productId;
    private String category;
    private double amount;
    private long timestamp;
    // getters/setters...
}

KAFKA STREAMS AGGREGATION :
// SalesAggregator.java
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

public class SalesAggregator {
    
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "sales-aggregator");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, 
                  Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, 
                  Serdes.String().getClass());
        
        StreamsBuilder builder = new StreamsBuilder();
        
        // Source
        KStream<String, String> sales = builder.stream("sales");
        
        // Parse JSON
        KStream<String, Sale> parsedSales = sales.mapValues(json -> 
            new ObjectMapper().readValue(json, Sale.class)
        );
        
        // Agrégation par catégorie (fenêtre 1 minute)
        KTable<Windowed<String>, SalesStats> categoryStats = parsedSales
            .groupBy((key, sale) -> sale.getCategory())
            .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
            .aggregate(
                SalesStats::new,  // Initializer
                (category, sale, stats) -> {  // Aggregator
                    stats.incrementCount();
                    stats.addAmount(sale.getAmount());
                    return stats;
                },
                Materialized.with(Serdes.String(), salesStatsSerde)
            );
        
        // Output
        categoryStats
            .toStream()
            .map((windowedKey, stats) -> {
                String key = windowedKey.key() + "@" + 
                    windowedKey.window().start();
                String value = new ObjectMapper().writeValueAsString(stats);
                return KeyValue.pair(key, value);
            })
            .to("sales-stats-1min");
        
        // Démarrer
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();
        
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }
}

// SalesStats.java
public class SalesStats {
    private int count = 0;
    private double totalAmount = 0.0;
    private double avgAmount = 0.0;
    
    public void incrementCount() { count++; }
    public void addAmount(double amount) {
        totalAmount += amount;
        avgAmount = totalAmount / count;
    }
    // getters...
}

WEBSOCKET SERVER (Node.js) :
// server.js
const { Kafka } = require('kafkajs');
const WebSocket = require('ws');

const kafka = new Kafka({
  clientId: 'dashboard-server',
  brokers: ['localhost:9092']
});

const consumer = kafka.consumer({ groupId: 'dashboard-group' });

const wss = new WebSocket.Server({ port: 8080 });

// Clients connectés
const clients = new Set();

wss.on('connection', (ws) => {
  clients.add(ws);
  console.log('Client connected');
  
  ws.on('close', () => {
    clients.delete(ws);
    console.log('Client disconnected');
  });
});

// Consumer Kafka
async function run() {
  await consumer.connect();
  await consumer.subscribe({ topic: 'sales-stats-1min' });
  
  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      const stats = JSON.parse(message.value.toString());
      
      // Broadcast à tous les clients
      clients.forEach(client => {
        if (client.readyState === WebSocket.OPEN) {
          client.send(JSON.stringify(stats));
        }
      });
    },
  });
}

run().catch(console.error);

DASHBOARD (React) :
// Dashboard.jsx
import React, { useState, useEffect } from 'react';
import { Line } from 'react-chartjs-2';

function Dashboard() {
  const [salesData, setSalesData] = useState({
    labels: [],
    datasets: []
  });
  
  useEffect(() => {
    const ws = new WebSocket('ws://localhost:8080');
    
    ws.onmessage = (event) => {
      const stats = JSON.parse(event.data);
      
      // Update chart
      setSalesData(prevData => ({
        labels: [...prevData.labels, new Date().toLocaleTimeString()],
        datasets: [{
          label: stats.category,
          data: [...(prevData.datasets[0]?.data || []), stats.totalAmount],
          borderColor: 'rgb(75, 192, 192)',
          tension: 0.1
        }]
      }));
    };
    
    return () => ws.close();
  }, []);
  
  return (
    <div>
      <h1>Real-Time Sales Dashboard</h1>
      <Line data={salesData} />
    </div>
  );
}

export default Dashboard;


EXEMPLE #3 : CDC PIPELINE (MySQL -> Elasticsearch)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF :
Synchroniser MySQL vers Elasticsearch en temps réel

ARCHITECTURE :
MySQL -> Debezium -> Kafka -> Kafka Connect Sink -> Elasticsearch

ÉTAPE 1 : ACTIVER BINLOG MYSQL
# /etc/mysql/my.cnf
[mysqld]
server-id = 1
log_bin = mysql-bin
binlog_format = ROW
binlog_row_image = FULL
expire_logs_days = 7

ÉTAPE 2 : DEBEZIUM SOURCE CONNECTOR
curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{
  "name": "mysql-debezium-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "mysql",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "secret",
    "database.server.id": "184054",
    "database.server.name": "mydb",
    "table.include.list": "mydb.products,mydb.orders",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.mydb",
    "snapshot.mode": "initial",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite"
  }
}'

ÉTAPE 3 : ELASTICSEARCH SINK CONNECTOR
curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{
  "name": "elasticsearch-sink-connector",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "tasks.max": "2",
    "topics": "mydb.products,mydb.orders",
    "connection.url": "http://elasticsearch:9200",
    "type.name": "_doc",
    "key.ignore": "false",
    "schema.ignore": "true",
    "behavior.on.null.values": "delete",
    "transforms": "route",
    "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.route.regex": "mydb\\.(.*)",
    "transforms.route.replacement": "$1"
  }
}'

VÉRIFICATION :
# Topic créé automatiquement par Debezium
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list | grep mydb

# Consommer quelques messages
bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:9092 \
  --topic mydb.products \
  --from-beginning \
  --max-messages 5

# Vérifier dans Elasticsearch
curl -X GET "http://localhost:9200/products/_search?pretty"

RÉSULTATS :
[OK] Sync temps réel (< 1 seconde)
[OK] Capture INSERT, UPDATE, DELETE
[OK] Schema evolution automatique
[OK] Full audit trail


EXEMPLE #4 : FRAUD DETECTION REAL-TIME
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

OBJECTIF :
Détecter transactions frauduleuses en temps réel

RÈGLES DE DÉTECTION :
-> Montant > 10,000€
-> 3+ transactions en 5 minutes
-> Géolocalisation incohérente (ex: Paris puis Tokyo en 1h)

KAFKA STREAMS IMPLEMENTATION :
// FraudDetector.java
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.*;

public class FraudDetector {
    
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "fraud-detector");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
                  StreamsConfig.EXACTLY_ONCE_V2);
        
        StreamsBuilder builder = new StreamsBuilder();
        
        // Source : transactions
        KStream<String, Transaction> transactions = builder
            .stream("transactions", Consumed.with(Serdes.String(), transactionSerde));
        
        // Règle 1 : Montant élevé
        KStream<String, Transaction> highAmount = transactions
            .filter((key, tx) -> tx.getAmount() > 10000);
        
        // Règle 2 : Fréquence (3+ en 5 min)
        KTable<Windowed<String>, Long> txCount = transactions
            .groupByKey()
            .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
            .count();
        
        KStream<String, Long> highFrequency = txCount
            .toStream()
            .filter((windowed, count) -> count >= 3)
            .selectKey((windowed, count) -> windowed.key());
        
        // Règle 3 : Géolocalisation (custom processor)
        KStream<String, Transaction> geoAnomaly = transactions
            .transformValues(() -> new GeoAnomalyDetector(), 
                           Named.as("geo-detector"),
                           "user-locations-store");
        
        // Merger toutes les alertes
        KStream<String, FraudAlert> alerts = highAmount
            .mapValues(tx -> new FraudAlert(tx, "HIGH_AMOUNT"))
            .merge(highFrequency.join(
                transactions,
                (count, tx) -> new FraudAlert(tx, "HIGH_FREQUENCY")
            ))
            .merge(geoAnomaly
                .filter((key, tx) -> tx.isGeoAnomaly())
                .mapValues(tx -> new FraudAlert(tx, "GEO_ANOMALY"))
            );
        
        // Output
        alerts.to("fraud-alerts");
        
        // Démarrer
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();
    }
}

// GeoAnomalyDetector.java
public class GeoAnomalyDetector implements ValueTransformer<Transaction, Transaction> {
    
    private KeyValueStore<String, UserLocation> locationStore;
    
    @Override
    public void init(ProcessorContext context) {
        this.locationStore = context.getStateStore("user-locations-store");
    }
    
    @Override
    public Transaction transform(Transaction tx) {
        String userId = tx.getUserId();
        UserLocation lastLocation = locationStore.get(userId);
        
        if (lastLocation != null) {
            double distance = calculateDistance(
                lastLocation.getLatitude(), lastLocation.getLongitude(),
                tx.getLatitude(), tx.getLongitude()
            );
            
            long timeDiff = tx.getTimestamp() - lastLocation.getTimestamp();
            double maxPossibleDistance = (timeDiff / 3600000.0) * 900;  // 900km/h
            
            if (distance > maxPossibleDistance) {
                tx.setGeoAnomaly(true);
            }
        }
        
        // Mettre à jour location
        locationStore.put(userId, new UserLocation(
            tx.getLatitude(), tx.getLongitude(), tx.getTimestamp()
        ));
        
        return tx;
    }
    
    private double calculateDistance(double lat1, double lon1, double lat2, double lon2) {
        // Formule Haversine
        double R = 6371;  // Rayon Terre en km
        double dLat = Math.toRadians(lat2 - lat1);
        double dLon = Math.toRadians(lon2 - lon1);
        double a = Math.sin(dLat/2) * Math.sin(dLat/2) +
                   Math.cos(Math.toRadians(lat1)) * Math.cos(Math.toRadians(lat2)) *
                   Math.sin(dLon/2) * Math.sin(dLon/2);
        double c = 2 * Math.atan2(Math.sqrt(a), Math.sqrt(1-a));
        return R * c;
    }
}

ALERT CONSUMER :
// AlertConsumer.java
public class AlertConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "alert-processor");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
        consumer.subscribe(Collections.singletonList("fraud-alerts"));
        
        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
            
            for (ConsumerRecord<String, String> record : records) {
                FraudAlert alert = parseAlert(record.value());
                
                // Actions
                blockCard(alert.getUserId());
                sendSMS(alert.getUserId(), "Suspicious transaction detected");
                createCase(alert);
                
                System.out.println("FRAUD ALERT: " + alert);
            }
        }
    }
}

RÉSULTATS :
[OK] Détection en < 100ms
[OK] Exactly-once (pas de faux négatifs)
[OK] State stores pour historique utilisateur
[OK] 99.9% précision


RÉSUMÉ DES BONNES PRATIQUES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRODUCTION CHECKLIST :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

BROKER :
[OK] replication.factor >= 3
[OK] min.insync.replicas = 2
[OK] unclean.leader.election.enable = false
[OK] auto.create.topics.enable = false
[OK] log.retention configured appropriately
[OK] Monitoring (Prometheus + Grafana)
[OK] Alerting (Under Replicated Partitions, Offline Partitions)

PRODUCER :
[OK] enable.idempotence = true
[OK] acks = all (données critiques)
[OK] compression.type = zstd
[OK] Batching (linger.ms + batch.size)
[OK] Error handling (callbacks)
[OK] Monitoring (send rate, error rate)

CONSUMER :
[OK] enable.auto.commit = false
[OK] Manual commit after processing
[OK] isolation.level = read_committed (exactly-once)
[OK] max.poll.interval.ms configured
[OK] Idempotent processing
[OK] Monitoring (consumer lag)

KAFKA STREAMS :
[OK] EXACTLY_ONCE_V2
[OK] State stores avec changelog
[OK] Exception handlers
[OK] Testing (TopologyTestDriver)

KAFKA CONNECT :
[OK] Distributed mode
[OK] Error handling (DLQ)
[OK] Avro + Schema Registry
[OK] Monitoring

GÉNÉRAL :
[OK] Capacity planning
[OK] Disaster recovery plan
[OK] Documentation
[OK] Testing (performance, failover)
[OK] Security (TLS, SASL)
"""


━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
                           FIN DU GUIDE KAFKA
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Ce guide couvre Apache Kafka de A à Z, des fondamentaux aux cas d'usage avancés en production. Pour aller plus loin :

[DOCS] RESSOURCES :
-> Documentation officielle : https://kafka.apache.org/documentation/
-> Confluent Developer : https://developer.confluent.io/
-> Kafka: The Definitive Guide (O'Reilly)
-> Designing Data-Intensive Applications (Martin Kleppmann)

[COURS] CERTIFICATIONS :
-> Confluent Certified Developer for Apache Kafka (CCDAK)
-> Confluent Certified Administrator for Apache Kafka (CCAAK)

[SPEECH_BALLOON] COMMUNAUTÉ :
-> Kafka Users mailing list
-> Confluent Community Slack
-> Stack Overflow (tag: apache-kafka)

Bon développement avec Kafka ! [RAPIDE]