control-systems-and-automation
Comment utiliser Kafka pour construire des applications robustes d'événements
Table of Contents
Comprendre Apache Kafka et son rôle dans l'architecture animée par des événements
Apache Kafka est une plateforme de streaming d'événements distribuée capable de gérer des trillions d'événements par jour. Initialement développée à LinkedIn, Kafka est devenue l'épine dorsale des architectures modernes axées sur les événements, permettant aux applications de publier, stocker, traiter et réagir aux flux de données en temps réel. Sa capacité à combiner un débit élevé, une tolérance aux défauts et une évolutivité horizontale en fait un choix idéal pour construire des systèmes robustes et de qualité de production axés sur les événements.
Ce qui distingue Kafka des files d'attente traditionnelles est sa conception centrale comme un journal de commit distribué. Au lieu de supprimer les messages après consommation, Kafka les conserve pour une période configurable (ou pour toujours), permettant à plusieurs consommateurs de rejouer ou de reproduire des événements. Ce découplage des producteurs et des consommateurs signifie que chaque côté peut s'étendre indépendamment, et les défaillances dans une partie du système ne cascade pas. Pour les applications axées sur les événements, ce choix architectural se traduit directement en robustesse : vous pouvez ajouter de nouveaux consommateurs sans perturber ceux existants, et vous pouvez récupérer des échecs en reluant simplement un offset connu.
Les composantes de base de Kafka : une plongée plus profonde
Pour construire des applications robustes basées sur l'événement avec Kafka, vous devez d'abord saisir ses éléments fondamentaux. Chaque composant joue un rôle essentiel dans la performance et la fiabilité de la plateforme :
- Les sujets sont des canaux logiques vers lesquels les enregistrements sont publiés. Un sujet peut avoir un certain nombre de partitions, et la stratégie de partitionnement détermine comment les données sont distribuées entre courtiers.
- Les partitions sont l'unité de parallélisme et de commande. Dans une partition, les enregistrements sont strictement commandés par offset. Les producteurs peuvent choisir une clé de partition (p. ex., ID utilisateur) pour s'assurer que tous les événements pour la même clé vont à la même partition, en préservant l'ordre de cette entité.
- Les producteurs publient des enregistrements sur des sujets. Ils peuvent configurer les reconnaissances (paquets) pour équilibrer la vitesse par rapport à la durabilité:
- – pas de reconnaissance, plus rapide mais risque de perte de données.
- – le leader reconnaît, bon équilibre.
- – toutes les répliques in-sync reconnaissent, la durabilité la plus forte.
- Les consommateurs lisent les enregistrements des partitions. Ils appartiennent à un groupe de consommateurs, ce qui permet l'équilibrage de la charge : chaque partition est attribuée à un seul consommateur du groupe. Si un consommateur échoue, les partitions sont rééquilibrées aux autres membres, garantissant qu'aucune donnée ne passe sans traitement.
- Les braqueurs sont des serveurs Kafka qui stockent des données et servent les demandes des clients. Un cluster Kafka se compose généralement de plusieurs courtiers. Chaque partition est reproduite sur un nombre configurable de courtiers (facteur de réplication) pour fournir une tolérance de défaillance.
Comprendre comment ces composants interagissent est crucial pour concevoir un déploiement Kafka qui répond aux exigences de votre application pour le débit, la latence, la durabilité et la cohérence.
Mise en place de Kafka pour la production-réapprovisionnement en flux d'événements
Une configuration de développement avec un seul courtier est idéale pour l'apprentissage, mais une application robuste axée sur l'événement exige une configuration de production. Voici les étapes clés et les considérations:
Taille des grappes et configuration des courtiers
Commencez par au moins trois courtiers pour assurer le quorum pour l'élection du chef et permettre la maintenance sans temps d'arrêt. Configurez le facteur de réplication à 3 pour les sujets critiques. Définissez à 2 pour garantir qu'au moins deux répliques reconnaissent les écrits lors de l'utilisation .
Thème Conception et stratégie de partage
Une bonne règle consiste à commencer par 10 à 50 partitions par sujet, selon le débit prévu. Chaque partition est essentiellement un fichier, de sorte que trop de partitions peuvent conduire à gérer les fichiers en hauteur et augmenter la charge de Zookeeper. Envisagez d'utiliser les lignes directrices pour le dimensionnement de partitions de confluents pour votre charge de travail spécifique. Utilisez des touches de partition significatives (p. ex., identifiant de commande, identifiant du client) pour préserver l'ordre dans le flux d'événements de l'entité.
Intégration au registre du schéma confluent
Pour maintenir la compatibilité des données à mesure que vos schémas d'événements évoluent, intégrez le registre du schéma confluent. Ce service stocke les définitions d'Avro, de Protobuf ou de JSON Schema et applique les règles de compatibilité (en arrière, en avant, en plein). Les producteurs et les consommateurs se réfèrent à l'ID du schéma plutôt qu'à l'intégration de schémas complets, réduisant ainsi les frais généraux du réseau.
Mise en oeuvre des pratiques exemplaires des producteurs et des consommateurs
Kafka offre de riches bibliothèques clientes pour Java, Python, Go, .NET et bien d'autres langues. Les exemples suivants utilisent Java, mais les modèles s'appliquent universellement.
Créer un producteur fiable
Un producteur robuste devrait gérer les prélèvements, l'idempotence et la sémantique transactionnelle:
- Activer l'idempotence en paramétrant . Cela empêche les enregistrements en double en cas de retraits, assurant exactement une sémantique pour les écritures à partition unique.
- Régler à une valeur élevée (p. ex. ) et configurer à des entrées liées.
- Utilisez asynchrone des envois avec un callback pour gérer les échecs gracieusement : enregistrez l'erreur, l'alerte ou la route vers un sujet de lettre morte.
- Choisissez un partitionneur qui distribue uniformément la charge. Le partitionneur collant par défaut améliore l'efficacité de lotage.
Exemple d'extrait de code (pseudocode):
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
props.put("enable.idempotence", true);
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("orders", orderKey, orderBytes), (metadata, exception) -> {
if (exception != null) {
// handle exception – log, alert, send to DLT
}
});
Créer un consommateur résilient
Les consommateurs doivent gérer gracieusement le rééquilibrage, gérer les compensations et traiter de façon idéo-pétente :
- Définir et commit manuellement des offsets après le traitement d'un lot. Cela empêche la perte de données si le consommateur s'écrase avant de commettre.
- Utilisez pour contrôler la taille des lots et éviter de traiter trop de dossiers avant de les engager.
- Mettre en place un auditeur rééquilibre pour stocker les compensations avant la révocation de partition et chercher à stocker les compensations lors de l'assignation.
- Faire un idémpotent de traitement de sorte que les duplications du retraitement ne causent pas d'effets secondaires. Par exemple, dédupliquer par ID d'événement ou utiliser une base de données upsert.
Pour les données à haut débit, envisager d'utiliser une boucle poll [ qui traite les enregistrements en parallèle à l'aide d'un pool de fils, mais s'assurer que les commits offset ne se produisent qu'après le traitement de tous les enregistrements d'un lot. La documentation de consommation d'Apache Kafka offre une plongée profonde sur ces mécanismes.
Traitement avancé des événements avec Kafka Streams et KSQL
Au-delà de la simple production/consommation, Kafka fournit des capacités de traitement de flux de première classe.
Fuseaux de Kafka
Kafka Streams est une bibliothèque cliente pour la construction d'applications de streaming d'état. Elle fonctionne comme une application standard (pas de cluster séparé) et exploite les propres sujets de Kafka pour les magasins d'état et les changelogs.
- Exactement une fois sémantique pour les opérations d'état (joins, agrégations).
- Support natif pour la fenêtre (montage, saut, fenêtres de session).
- API et DSL du processeur (p. ex. ).
Par exemple, vous pouvez calculer un total de commandes en cours d'exécution par client en créant une KTable à partir d'un sujet de commande et en utilisant l'opérateur . Kafka Streams gère automatiquement le magasin d'état et changelog, rendant votre application automatiquement résilient aux défaillances – si un noeud s'écrase, l'état est reconstruit à partir du sujet de changement.
KSQL (Kafka SQL)
KSQL est le moteur SQL de diffusion de Kafka. Il vous permet d'exécuter des requêtes SQL sur les données de diffusion sans écrire de code Java. Utilisez-le pour l'analyse ad-hoc, le prototypage ou l'ETL simple. Par exemple:
CREATE STREAM orders WITH (KAFKA_TOPIC='orders', VALUE_FORMAT='JSON');
CREATE TABLE high_value_orders AS
SELECT customer_id, COUNT(*) AS order_count, SUM(amount) AS total
FROM orders WINDOW TUMBLING (SIZE 1 HOUR)
WHERE amount > 1000
GROUP BY customer_id;
KSQL est particulièrement utile pour les équipes d'ingénierie des données qui veulent construire des transformations animées par des événements rapidement.
Meilleures pratiques pour construire des systèmes de production robustes
Une application résiliente axée sur les événements va au-delà de la simple écriture des producteurs et des consommateurs. Elle nécessite une approche holistique de la conception, des opérations et de la surveillance.
Gestion des erreurs et requêtes en lettres mortes
Même avec des consommateurs robustes, certains documents ne seront pas traités (p. ex., JSON malformé, pannes transitoires en aval). Mettre en place un modèle où le consommateur prend des exceptions, enregistre le document original et le publie sur un sujet sans objet (p. ex. ). Un processus distinct peut ensuite rejouer ces documents après enquête.
Garantie de la sémantique une fois
Pour les applications où les duplicata sont inacceptables (p. ex., transactions financières), utiliser la sémantique exacte de Kafka (EOS) pour les producteurs et les consommateurs. Du côté des producteurs, comme mentionné, ne garantit pas de duplicata au cours d'une session. Du côté des consommateurs, utiliser l'API transactionnelle pour écrire les enregistrements de sortie et les offsets atomiquement.
Surveillance et observation
Kafka expose de nombreuses mesures via JMX.
- partitions sous-réplicables:[ Indique un problème avec la réplication.
- Dégaiement des consommateurs:[ Différence entre le dernier décalage et le décalage engagé du consommateur.
- Demander la latence: Temps de production ou de consommation.
Utilisez des outils comme Prométhée avec l'exportateur Kafka JMX pour collecter des métriques et installer des tableaux de bord dans Grafana. En outre, activez l'analyseur de log intégré de Kafka (par exemple ) pour le débogage.
Pratiques exemplaires en matière de sécurité
Protégez vos données en transit et au repos :
- Authentification: Utilisez SASL/SCRAM ou SASL/SSL pour l'authentification des clients.
- Autorisation:[ Définir des ACL pour contrôler quels utilisateurs peuvent lire/écrire sur des sujets.
- Encryptage:[ Activer TLS/SSL pour la communication client-courtier et courtier-courtier.
- Politiques de réseau:[ Utilisez des pare-feu et des VPC pour restreindre l'accès aux courtiers.
Pour un guide détaillé, consultez le Documentation de sécurité [.
Échelle et réglage
Avec l'augmentation du volume de votre événement, vous devrez peut-être ajuster le nombre de partitions, augmenter le facteur de réplication ou ajouter des courtiers. Planifiez la capacité en surveillant l'utilisation du disque, les E/S réseau et le processeur. Utilisez l'outil de Kafka pour rééquilibrer les données entre les nouveaux courtiers. Pour les scénarios à haut débit, accordez les tailles de lots (, ) pour les producteurs et récupérez les tailles pour les consommateurs.
Cas et modèles d'utilisations dans le monde réel
Pour illustrer comment ces concepts se réunissent, considérez une plateforme typique de commerce électronique qui utilise Kafka comme système nerveux central :
- Order Service publie des événements « OrderPlaced » sur un sujet .
- Le service d'inventaire consomme ces événements pour réserver des stocks, puis publie « InventaireRéservé» ou «OutOfStock».
- Le Service de paiement consomme les événements et traite les paiements « InventaireRéservé», publiant «Paiement complété».
- Notification Service consomme "Paiement Complété" et envoie des confirmations de courriel/SMS.
- Le service d'analyse consomme tous les événements de commande pour construire un tableau de bord en temps réel.
- Une application Kafka Streams rejoint les flux d'événements pour détecter les tendances de fraude (par exemple, trop de commandes de la même IP en peu de temps).
Dans cette architecture, chaque service s'évalue de façon indépendante. Si le service de notification est en panne pour maintenance, les événements restent à Kafka et sont traités plus tard. Si le service de paiement échoue après avoir commis, l'événement PaymentComplete assure la récupération idempotent. L'utilisation d'un registre de schéma garantit que lorsque le service de commande ajoute un nouveau champ (p. ex., « code de retrait »), les services en aval ne sont pas immédiatement brisés.
Un autre motif courant est le modèle Event Sourcing, où la source principale de vérité est le flux d'événement lui-même. Le journal de Kafka en appendice sert de magasin d'événements. Les services d'État reconstruisent leur état en rejouant les événements dès le début (ou à partir d'un instantané). Ce modèle fournit une piste d'audit complète et la possibilité de corriger rétroactivement les bogues en rejouant les événements corrigés.
Comparaison avec d'autres technologies d'événements
Bien que Kafka soit puissante, ce n'est pas la seule solution. Comprendre quand l'utiliser par rapport aux alternatives vous aidera à faire le bon choix architectural :
- RabbitMQ excelle dans la messagerie point à point à faible latence avec routage complexe (échanges, fixations). Il est plus léger pour les déploiements plus petits mais manque de garanties de durabilité de Kafka et de capacité de rejouer. Utilisez RabbitMQ lorsque vous avez besoin de livraison garantie à un seul consommateur avec des frais généraux faibles.
- Amazon Kinesis est un service de streaming géré semblable à Kafka, mais il élimine les frais généraux opérationnels. Cependant, il peut avoir un coût plus élevé à l'échelle et moins de flexibilité dans le réglage. Kafka offre plus de contrôle et sur site options de déploiement.
- Apache Pulsar fournit un stockage à plusieurs niveaux et des logements multiples nativement, mais a une communauté plus petite et moins d'outils écosystémiques. La maturité de Kafka, une communauté massive et de vastes bibliothèques clientes en font souvent le choix plus sûr pour les systèmes à grande échelle axés sur les événements.
Finalement, Kafka est le meilleur pour les applications qui nécessitent des flux d'événements commandés, durables, rejouables avec un débit élevé et faible latence, en particulier lors de l'intégration de plusieurs microservices ou de la construction d'un lac de données.
Conclusion
Pour construire des applications robustes basées sur les événements avec Apache Kafka, il faut plus que comprendre son API – il faut une compréhension approfondie de son architecture, une configuration soignée pour la production et le respect des meilleures pratiques de gestion des erreurs, de surveillance et de sécurité. En exploitant les composants de base de Kafka (sujets, partitions, producteurs, consommateurs, courtiers) et les capacités avancées comme Kafka Streams et le Schema Registry, vous pouvez créer des systèmes résilients en cas de défaillance, évolutives à des charges élevées et durables au fil du temps.
Commencez par modéliser vos événements avec soin, concevoir vos sujets avec une croissance future en tête, et toujours planifier pour l'inattendu : partitions réseau, brokers s'écraser, et changements de schéma. Avec Kafka, vous gagnez la capacité de découpler les services, de permettre le flux de données en temps réel, et de construire des applications qui non seulement survivent mais prospèrent face à la complexité. Pour plus de lecture, explorez la documentation Apache Kafka et la bibliothèque de ressources confluente pour des guides approfondis et des architectures de référence.