Architecture Kafka Connect expliquée : Source Connectors, Sink Connectors & CDC
Les architectures modernes d’entreprise stockent rarement toutes leurs données dans une seule application.
Les informations clients peuvent se trouver dans PostgreSQL, les commandes dans MySQL, les événements dans Apache Kafka, les données analytiques dans un Data Warehouse et les données de recherche dans Elasticsearch.
Déplacer ces données de manière fiable entre tous ces systèmes peut rapidement devenir complexe.
C’est précisément le problème que Apache Kafka Connect permet de résoudre.
Kafka Connect fournit un framework standardisé permettant de transférer les données vers et depuis Apache Kafka, sans obliger chaque équipe de développement à créer et maintenir ses propres producteurs et consommateurs Kafka pour les scénarios d’intégration courants.
Dans ce guide, nous allons découvrir :
- l’architecture Kafka Connect ;
- les Source Connectors ;
- les Sink Connectors ;
- les Connect Workers ;
- les Connectors et Tasks ;
- les Converters ;
- les Single Message Transforms (SMT) ;
- la gestion des offsets ;
- les modes Standalone et Distributed ;
- le Change Data Capture (CDC) ;
- Debezium avec Kafka Connect ;
- la gestion des erreurs ;
- la scalabilité et la tolérance aux pannes ;
- les bonnes pratiques pour la Production.
1. Qu’est-ce que Kafka Connect ?
Kafka Connect est un framework d’intégration permettant de transférer des données entre Apache Kafka et des systèmes externes.
Au lieu de développer du code personnalisé pour chaque base de données, système de fichiers, plateforme de recherche ou Data Warehouse, vous configurez un connecteur adapté.
L’architecture de base est simple :
Système externe ↓ Source Connector ↓ Kafka Connect ↓ Kafka Topics ↓ Kafka Connect ↓ Sink Connector ↓ Système cible
Kafka Connect définit deux grandes catégories de connecteurs.
Source Connector
Il déplace les données :
Système externe → Kafka
Sink Connector
Il déplace les données :
Kafka → Système externe
2. Pourquoi utiliser Kafka Connect?
Imaginons qu’une entreprise doive réaliser les intégrations suivantes :
PostgreSQL → Kafka MySQL → Kafka Kafka → Elasticsearch Kafka → Data Warehouse Kafka → Object Storage
Une première solution consisterait à développer une application spécifique pour chaque intégration.
Il faudrait alors gérer :
- les Kafka Producers ;
- les Kafka Consumers ;
- la sérialisation ;
- les offsets ;
- les retries ;
- les erreurs ;
- la scalabilité ;
- le monitoring ;
- le déploiement ;
- la configuration.
Kafka Connect fournit un framework commun prenant en charge une grande partie de cette infrastructure.
Les équipes peuvent ainsi se concentrer davantage sur les données à transférer et leur destination plutôt que de reconstruire constamment la couche technique d’intégration.
3. Architecture de Kafka Connect
Un environnement Kafka Connect comprend plusieurs composants importants :
Système externe ↓ Connector ↓ Tasks ↓ Worker ↓ Converters / SMT ↓ Kafka Cluster
Voici leurs responsabilités principales :
| Composant | Responsabilité |
|---|---|
| Connector | Définit et coordonne l’intégration |
| Task | Exécute le transfert des données |
| Worker | Exécute les Connectors et Tasks |
| Converter | Convertit la représentation des données |
| SMT | Transforme légèrement chaque record |
| Kafka Topics | Stockent les événements |
| Offsets | Suivent la progression |
| REST API | Permet de gérer les Connectors |
Le Connector définit le travail à effectuer, tandis que les Tasks exécutent concrètement les opérations de transfert.
4. Qu’est-ce qu’un Kafka Connect Worker ?
Un Worker est un processus Kafka Connect en cours d’exécution.
Il peut exécuter :
- des Connectors ;
- des Tasks ;
- des Converters ;
- des transformations.
Par exemple :
Kafka Connect Cluster | +---- Worker 1 | | | +---- Task A | +---- Task B | +---- Worker 2 | | | +---- Task C | +---- Worker 3 | +---- Task D
En mode distribué, plusieurs Workers fonctionnent ensemble pour constituer un Kafka Connect Cluster.
Cette architecture permet de répartir les Tasks et d’améliorer la capacité et la disponibilité du système.
5. Connectors vs Tasks
Cette distinction est essentielle pour comprendre Kafka Connect.
Connector
Le Connector gère la configuration de l’intégration et détermine comment le travail peut être réparti.
Task
La Task effectue réellement le transfert des données.
Par exemple :
JDBC Source Connector | +---- Task 1 +---- Task 2 +---- Task 3
Une configuration peut contenir :
{ "tasks.max": "3" }
Cela signifie que le Connector peut utiliser jusqu’à trois Tasks.
Cela ne signifie pas nécessairement que trois Tasks seront toujours créées. Le nombre réel dépend également de la capacité du Connector et de la source à paralléliser le travail.
6. Qu’est-ce qu’un Source Connector ?
Un Source Connector importe les données d’un système externe vers Kafka.
Le sens du flux est :
Système externe ↓ Source Connector ↓ Kafka Topic
Les sources peuvent notamment être :
- des bases de données relationnelles ;
- des fichiers ;
- des applications SaaS ;
- des systèmes de messagerie ;
- du stockage Cloud ;
- des logs applicatifs ;
- des plateformes CDC.
Exemple :
PostgreSQL ↓ Debezium PostgreSQL Connector ↓ Kafka Connect ↓ Kafka Topic
7. Exemple de Source Connector
Une configuration conceptuelle peut ressembler à :
{ "name": "customer-source", "config": { "connector.class": "com.example.CustomerSourceConnector", "tasks.max": "2", "topic": "customers" } }
Parmi les propriétés courantes :
name connector.class tasks.max key.converter value.converter
Les autres paramètres dépendent du Connector utilisé.
8. Qu’est-ce qu’un Sink Connector ?
Un Sink Connector transfère les données depuis les topics Kafka vers un système externe.
Kafka Topic ↓ Sink Connector ↓ Système cible
Les destinations courantes peuvent être :
- bases de données ;
- Elasticsearch ;
- Data Warehouses ;
- Object Storage ;
- plateformes analytiques ;
- systèmes de fichiers ;
- applications externes.
Par exemple :
Kafka Topic ↓ Sink Connector ↓ Data Warehouse
9. Exemple de Sink Connector
Une configuration simplifiée peut être :
{ "name": "customer-sink", "config": { "connector.class": "com.example.CustomerSinkConnector", "tasks.max": "2", "topics": "customers" } }
Le Sink Connector s’abonne généralement à un ou plusieurs topics Kafka puis écrit les records dans le système cible.
10. Source Connector vs Sink Connector
La façon la plus simple de retenir la différence est de se placer du point de vue de Kafka.
| Source Connector | Sink Connector |
|---|---|
| Système externe → Kafka | Kafka → Système externe |
| Produit des records dans Kafka | Consomme les records Kafka |
| Lit les données de la source | Écrit dans la destination |
| Peut être utilisé pour CDC | Utilisé pour les systèmes downstream |
| Exemple : DB → Kafka | Exemple : Kafka → Elasticsearch |
Retenez simplement :
SOURCE = données entrant dans Kafka
SINK = données sortant de Kafka
11. Pipeline Kafka Connect de bout en bout
Une intégration complète peut ressembler à :
PostgreSQL ↓ Debezium Source Connector ↓ Kafka Connect ↓ customer-events ↓ Kafka Cluster ↓ Sink Connector ↓ Analytics / Search / Data Warehouse
L’un des grands avantages de cette architecture est le découplage.
La base de données source n’a pas besoin de connaître tous les systèmes qui consommeront ses changements.
Plusieurs applications peuvent également exploiter indépendamment les mêmes événements Kafka.
12. Que sont les Kafka Connect Converters ?
Les Connectors transfèrent des records, mais il faut également gérer la représentation de ces données.
C’est le rôle des Converters.
Les formats courants comprennent :
JSON Avro Protobuf String ByteArray
La configuration peut notamment définir :
key.converter=... value.converter=...
Le Converter contrôle la manière dont les données Kafka Connect sont représentées lorsqu’elles sont écrites dans Kafka ou lues depuis Kafka.
Cela permet de séparer la logique d’intégration du format de sérialisation.
13. Que sont les Single Message Transforms (SMT) ?
Kafka Connect propose les Single Message Transforms, ou SMT.
Ils permettent d’effectuer de petites transformations sur chaque record traversant Kafka Connect.
Côté Source :
Source ↓ Connector ↓ SMT ↓ Converter ↓ Kafka
Côté Sink :
Kafka ↓ Converter ↓ SMT ↓ Sink Connector ↓ Destination
Les SMT peuvent notamment servir à :
- renommer des champs ;
- supprimer des champs ;
- ajouter des métadonnées ;
- modifier le routage ;
- modifier les noms de topics ;
- extraire certaines parties d’un record.
Les transformations complexes sont généralement mieux adaptées à Kafka Streams ou à une autre plateforme de stream processing.
14. Qu’est-ce que le Change Data Capture (CDC) ?
Le Change Data Capture, généralement appelé CDC, consiste à détecter les changements réalisés dans une base de données et à propager ces modifications vers d’autres systèmes.
Supposons que le statut d’un client passe de :
status = PENDING
à :
status = APPROVED
Au lieu d’interroger constamment toute la table, un système CDC peut détecter ce changement et produire un événement.
Transaction Database ↓ Transaction / Replication Log ↓ CDC Connector ↓ Kafka Connect ↓ Kafka Topic
Le CDC est particulièrement utile pour :
- les architectures event-driven ;
- les microservices ;
- l’analytics ;
- la synchronisation de caches ;
- l’indexation de recherche ;
- la réplication des données ;
- les pipelines d’audit ;
- les Data Warehouses.
15. Debezium + Kafka Connect
Debezium est largement utilisé pour implémenter le CDC avec Kafka Connect.
Les Connectors Debezium surveillent les changements dans les bases de données et publient ces changements sous forme d’événements.
Architecture typique :
MySQL / PostgreSQL ↓ Database Change Log ↓ Debezium Source Connector ↓ Kafka Connect ↓ Kafka Topics
Pour MySQL, les changements peuvent être capturés depuis le binlog.
Pour PostgreSQL, le CDC peut exploiter le mécanisme de logical replication / WAL.
Les événements sont ensuite envoyés vers Kafka, où ils peuvent être consommés par plusieurs applications et systèmes downstream.
16. Pourquoi utiliser le CDC basé sur les logs ?
Une intégration traditionnelle peut exécuter périodiquement :
SELECT * FROM customers WHERE updated_at > ?
Il s’agit d’une approche de polling.
Le CDC basé sur les logs fonctionne différemment.
Il exploite les informations de changement générées par la base de données.
Cela présente plusieurs avantages :
- faible latence ;
- capture des INSERT ;
- capture des UPDATE ;
- capture des DELETE ;
- réduction du polling ;
- moins de requêtes répétitives ;
- meilleure adaptation aux architectures event-driven.
C’est pourquoi le CDC est particulièrement intéressant pour les intégrations proches du temps réel.
17. Exemple d’événement CDC
Supposons que cet enregistrement :
Customer ID: 101 Name: John Status: PENDING
devienne :
Customer ID: 101 Name: John Status: APPROVED
Un événement CDC pourrait conceptuellement contenir :
{ "before": { "id": 101, "name": "John", "status": "PENDING" }, "after": { "id": 101, "name": "John", "status": "APPROVED" }, "op": "u" }
La structure précise dépend du Connector et de sa configuration.
Les consommateurs downstream peuvent alors réagir au changement sans interroger directement la base opérationnelle.
18. Snapshot initial + CDC continu
Une question importante est :
Que se passe-t-il avec les données présentes avant le démarrage du Connector ?
De nombreuses implémentations CDC peuvent utiliser un snapshot initial, puis continuer à capturer les changements.
ÉTAPE 1 Base existante ↓ Initial Snapshot ↓ Kafka ÉTAPE 2 INSERT / UPDATE / DELETE ↓ Database Log ↓ CDC ↓ Kafka
Cette combinaison permet d’obtenir l’état initial puis de poursuivre avec les changements en continu.
19. Gestion des offsets dans Kafka Connect
Kafka Connect doit savoir quelles données ont déjà été traitées.
Cette progression est suivie à l’aide des offsets.
Pour un Source Connector, l’offset peut représenter :
- une position dans un fichier ;
- une position dans un journal de base de données ;
- un numéro de séquence ;
- un timestamp ;
- une position propre au système source.
Par exemple :
Record 1001 ✓ Record 1002 ✓ Record 1003 ✓ Record 1004 ← position actuelle Record 1005
Après un redémarrage, Kafka Connect peut ainsi reprendre le traitement à partir d’une position connue au lieu de recommencer systématiquement depuis le début.
20. Kafka Connect Standalone vs Distributed
Kafka Connect prend en charge deux modes principaux.
Standalone Mode
Single Connect Process ↓ Connectors ↓ Tasks
Ce mode est intéressant pour :
- le développement ;
- les tests locaux ;
- les environnements simples.
Mais il repose sur un processus unique.
Distributed Mode
Kafka Connect Cluster | +-- Worker 1 +-- Worker 2 +-- Worker 3
Ce mode offre notamment :
- la scalabilité ;
- la répartition de la charge ;
- une meilleure tolérance aux pannes ;
- la gestion centralisée des Connectors.
Pour les déploiements d’entreprise nécessitant disponibilité et scalabilité, Distributed Mode est généralement le choix approprié.
21. Topics internes de Kafka Connect
En mode distribué, Kafka Connect conserve plusieurs informations internes dans Kafka.
Les topics typiques sont :
connect-configs connect-offsets connect-status
Ils servent respectivement à gérer :
Configurations
Les configurations des Connectors.
Offsets
La progression du traitement.
Status
L’état des Connectors et Tasks.
Ces topics constituent donc une partie importante de l’architecture d’un cluster Kafka Connect distribué.
22. Kafka Connect REST API
Kafka Connect fournit une REST API permettant d’administrer les Connectors.
Quelques opérations courantes :
GET /connectors
Lister les Connectors.
POST /connectors
Créer un Connector.
GET /connectors/{name}/status
Consulter son état.
PUT /connectors/{name}/config
Modifier sa configuration.
DELETE /connectors/{name}
Supprimer un Connector.
Cette API facilite également l’automatisation des déploiements et des opérations.
23. Gestion des erreurs et Dead Letter Queue
Les données de Production ne sont jamais parfaitement propres.
Un record peut par exemple contenir :
- du JSON invalide ;
- un schéma inattendu ;
- un type incorrect ;
- un champ obligatoire absent ;
- une erreur de sérialisation.
Les stratégies possibles peuvent inclure :
Retry Skip Log Dead Letter Queue Fail Connector
Une architecture peut par exemple être :
Kafka Topic ↓ Sink Connector ↓ Record valide ? ↙ ↘ Oui Non ↓ ↓ Target DLQ Topic
Une Dead Letter Queue (DLQ) permet d’isoler les records problématiques afin de les analyser sans nécessairement bloquer tout le pipeline.
24. Architecture Kafka Connect en Production
Une architecture d’entreprise peut ressembler à :
DATABASES MySQL / PostgreSQL / Oracle ↓ CDC Connectors ↓ ┌─────────────────────┐ │ Kafka Connect │ │ Worker 1 │ │ Worker 2 │ │ Worker 3 │ └─────────────────────┘ ↓ Kafka Cluster Broker 1 / 2 / 3 ↓ ┌─────────────┼─────────────┐ ↓ ↓ ↓ Search Sink Warehouse Sink Applications
En Production, il faut également prendre en compte :
- l’authentification ;
- TLS ;
- SASL ;
- les ACL ;
- la gestion des secrets ;
- le monitoring ;
- les retries ;
- les DLQ ;
- la gestion des versions ;
- la compatibilité des schémas ;
- la capacité ;
- la reprise après incident.
25. Kafka Connect vs Producers et Consumers Kafka
Kafka Connect ne remplace pas tous les Producers et Consumers Kafka.
Utilisez Kafka Connect lorsque votre besoin concerne principalement une intégration standard avec un système externe.
Par exemple :
Database → Kafka Kafka → Elasticsearch Kafka → Data Warehouse
Utilisez plutôt un Producer ou Consumer personnalisé lorsque l’application contient une logique métier importante.
Order Service ↓ Complex Business Rules ↓ Kafka Producer
Une architecture réelle peut parfaitement utiliser les deux approches.
26. Kafka Connect vs Kafka Streams
Kafka Connect et Kafka Streams répondent à des besoins différents.
| Kafka Connect | Kafka Streams |
|---|---|
| Intégration des données | Traitement des streams |
| Déplace les données | Transforme et traite les données |
| Source/Sink Connectors | API Java |
| Systèmes externes ↔ Kafka | Kafka ↔ traitement ↔ Kafka |
| Principalement configuration | Principalement code applicatif |
Par exemple :
PostgreSQL ↓ Kafka Connect ↓ orders ↓ Kafka Streams ↓ validated-orders ↓ Kafka Connect ↓ Analytics Platform
Ces technologies sont donc complémentaires.
27. Kafka Connect + CDC + Microservices
Le CDC est particulièrement intéressant lors de la modernisation d’applications existantes.
Imaginons une application legacy enregistrant directement les commandes dans PostgreSQL.
Une architecture peut être :
Legacy Application ↓ PostgreSQL ↓ Database Log ↓ Debezium ↓ Kafka Connect ↓ order-events ↓ Microservices
Les nouveaux microservices peuvent ainsi réagir aux changements de données sans nécessairement imposer immédiatement une refonte complète de l’application legacy.
28. Monitoring de Kafka Connect
Un environnement Kafka Connect de Production doit être correctement supervisé.
Surveillez notamment :
- disponibilité des Workers ;
- état des Connectors ;
- état des Tasks ;
- Tasks en échec ;
- lag ;
- débit ;
- taux de retry ;
- volume DLQ ;
- mémoire JVM ;
- CPU ;
- réseau ;
- erreurs des Connectors.
Un Connector affiché comme :
RUNNING
ne garantit pas à lui seul que tout le pipeline fonctionne correctement.
Le flux de données de bout en bout doit également être surveillé.
29. Bonnes pratiques Kafka Connect
Pour un environnement de Production :
- utilisez le mode distribué lorsque la haute disponibilité est nécessaire ;
- déployez plusieurs Workers ;
- surveillez séparément les Connectors et les Tasks ;
- sécurisez la REST API Kafka Connect ;
- choisissez correctement vos Converters ;
- configurez explicitement les retries et la gestion des erreurs ;
- surveillez les Dead Letter Queues ;
- versionnez les configurations des Connectors ;
- ne stockez pas les mots de passe directement dans les configurations ;
- testez les mises à niveau des Connectors avant la Production ;
- testez le CDC avec des volumes réalistes ;
- documentez les mappings entre sources, topics et destinations.
30. Questions fréquentes sur Kafka Connect
Qu’est-ce que Kafka Connect ?
Kafka Connect est un framework permettant de transférer des données entre Apache Kafka et des systèmes externes.
Qu’est-ce qu’un Source Connector ?
Un Source Connector transfère les données :
Système externe → Kafka
Qu’est-ce qu’un Sink Connector ?
Un Sink Connector transfère les données :
Kafka → Système externe
Qu’est-ce qu’un Worker Kafka Connect ?
Un Worker est un processus Kafka Connect qui exécute les Connectors et les Tasks.
Qu’est-ce qu’une Task ?
Une Task réalise concrètement le travail de transfert de données attribué par un Connector.
Qu’est-ce que le CDC ?
Le Change Data Capture détecte les changements dans une base de données et permet de les propager vers d’autres systèmes.
Qu’est-ce que Debezium ?
Debezium fournit des Connectors CDC permettant notamment de capturer les changements de bases de données et de les publier sous forme d’événements.
Qu’est-ce qu’un SMT ?
Un Single Message Transform effectue une transformation légère sur chaque record traversant Kafka Connect.
Standalone ou Distributed en Production ?
Pour les environnements d’entreprise nécessitant scalabilité et résilience, le Distributed Mode est généralement préférable.
31. Résumé de l’architecture Kafka Connect
Retenez cette architecture :
SOURCE Database / File / API ↓ Source Connector ↓ Source Tasks ↓ Kafka Connect Workers ↓ Converters + SMT ↓ Kafka Topics SINK Kafka Topics ↓ Converters + SMT ↓ Kafka Connect Workers ↓ Sink Tasks ↓ Sink Connector ↓ Database / Search / Warehouse
Et pour le CDC :
Database ↓ Transaction Log ↓ Debezium ↓ Kafka Connect ↓ Kafka Topics ↓ Systèmes downstream
Conclusion
Kafka Connect fournit une approche standardisée et scalable pour connecter Apache Kafka aux bases de données, systèmes de fichiers, plateformes de recherche, solutions analytiques, services Cloud et autres technologies externes.
L’architecture devient beaucoup plus simple lorsque l’on distingue clairement ses différents composants.
Les Source Connectors font entrer les données dans Kafka.
Les Sink Connectors font sortir les données de Kafka.
Les Connectors définissent et coordonnent l’intégration.
Les Tasks réalisent le transfert effectif.
Les Workers exécutent les Connectors et les Tasks.
Les Converters contrôlent la représentation des données.
Les SMT réalisent des transformations légères.
Enfin, l’association de Kafka Connect et Debezium constitue une architecture particulièrement puissante pour le Change Data Capture (CDC) et les systèmes event-driven.
Elle permet de transformer les changements des bases de données en événements Kafka pouvant être exploités presque en temps réel par des microservices, moteurs de recherche, Data Warehouses, systèmes analytiques et autres applications.
Articles recommandés
Microservices Event-Driven avec Kafka & Spring Boot
Pour comprendre comment les microservices Spring Boot exploitent les événements Kafka.
Sécurité Kafka : SSL, SASL, ACL et gouvernance
Pour sécuriser Kafka, Kafka Connect et les communications avec les brokers.
Architecture de déploiement Kafka sur Kubernetes : Scaling, HA & Monitoring
Pour passer de l’architecture logique à un environnement Kafka de Production.
Kafka Consumer Groups Explained
Pour approfondir le fonctionnement des Consumers, partitions et mécanismes de scaling.
Spring Boot + Kafka : Event-Driven Microservices
Pour relier les pipelines Kafka Connect aux applications métier.
🎥 Learn IT with Shikha sur YouTube
Vous préférez apprendre en vidéo ?
Découvrez des tutoriels pratiques sur Apache Kafka, Spring Boot, Microservices, Camunda, Alfresco, Java et l’architecture d’entreprise.
S'abonner à Learn IT with Shikha sur YouTube
Vous pouvez également intégrer votre vidéo existante :
Kafka Consumer Groups Explained — Partitions, Offsets & Rebalancing
📢 Besoin d’aide pour Java, workflows ou backend?
J’aide les équipes à concevoir des applications scalables, performantes et prêtes pour la production.
Services:
- Développement Java & Spring Boot
- Implémentation workflows (Camunda, Flowable – BPMN, DMN)
- Intégrations API & microservices
- ECM & gestion documentaire (Alfresco)
- Optimisation performance & résolution incidents
🔗 https://shikhanirankari.blogspot.com/p/professional-services.html
📩 Email: ishikhanirankari@gmail.com | info@realtechnologiesindia.com
🌐 https://realtechnologiesindia.com
✔ Disponible pour consultation rapide
✔ Réponse sous 24 heures
🎥 Learn IT with Shikha on YouTube
Prefer learning through videos? Watch practical tutorials on Kafka, Camunda, Alfresco, Java, Spring Boot, Microservices and Enterprise Architecture.▶ Subscribe to Learn IT with Shikha on YouTube
Comments
Post a Comment