
Dans mon travail, je suis souvent confrontĂ© Ă de nouvelles solutions techniques ou produits logiciels pour lesquels il existe peu d'informations sur Internet en langue française. Cet article vise Ă combler une telle lacune en prenant un exemple de ma pratique rĂ©cente, oĂč il a Ă©tĂ© nĂ©cessaire de configurer l'envoi d'Ă©vĂ©nements CDC Ă partir de deux bases de donnĂ©es populaires (PostgreSQL et MongoDB) vers un cluster Kafka Ă l'aide de Debezium. J'espĂšre que cet article de synthĂšse, rĂ©digĂ© suite Ă ce travail, sera utile Ă d'autres.
Qu'est-ce que Debezium et le CDC ?
est un représentant de la catégorie des logiciels CDC (), plus précisément, c'est un ensemble de connecteurs pour diverses bases de données compatibles avec le framework Apache Kafka Connect.
C'est sous licence Apache License v2.0 et sponsorisé par la société Red Hat. Le développement a débuté en 2016 et à ce jour, il soutient officiellement les bases de données suivantes : MySQL, PostgreSQL, MongoDB, SQL Server. Il existe également des connecteurs pour Cassandra et Oracle, mais actuellement, ils sont en statut « accÚs anticipé », et les nouvelles versions ne garantissent pas la compatibilité avec les anciennes.
ComparĂ© Ă l'approche traditionnelle (oĂč l'application lit directement les donnĂ©es de la base de donnĂ©es), les principaux avantages du CDC incluent la mise en Ćuvre d'un streaming des modifications de donnĂ©es au niveau des lignes avec une faible latence, une haute fiabilitĂ© et une disponibilitĂ© accrue. Ces deux derniers points sont rĂ©alisĂ©s grĂące Ă l'utilisation d'un cluster Kafka comme stockage des Ă©vĂ©nements CDC.
Parmi les autres avantages, on peut noter que la mĂȘme reprĂ©sentation est utilisĂ©e pour stocker les Ă©vĂ©nements, donc l'application finale n'a pas Ă se soucier des nuances d'exploitation de diffĂ©rentes bases de donnĂ©es.
Enfin, grùce à l'utilisation d'un courtier de messages, il devient possible d'évoluer horizontalement les applications qui surveillent les changements de données. L'impact sur la source de données est minimisé, car l'obtention des données ne se fait pas directement à partir de la base de données, mais à partir du cluster Kafka.
Concernant l'architecture de Debezium
L'utilisation de Debezium se résume à un schéma simple :
Base de donnĂ©es (comme source de donnĂ©es) â connecteur dans Kafka Connect â Apache Kafka â consommateur
à titre d'illustration, voici un schéma fourni sur le site du projet :

Cependant, ce schéma ne me plaßt pas beaucoup, car il donne l'impression qu'il n'est possible d'utiliser que le connecteur sink.
En rĂ©alitĂ©, la situation est diffĂ©rente : le remplissage de votre Data Lake (le dernier maillon du schĂ©ma ci-dessus) Ââ ce n'est pas le seul moyen d'utiliser Debezium. Les Ă©vĂ©nements envoyĂ©s Ă Apache Kafka peuvent ĂȘtre utilisĂ©s par vos applications pour rĂ©soudre diverses situations. Par exemple :
- la suppression des données obsolÚtes du cache ;
- l'envoi de notifications ;
- les mises Ă jour des index de recherche ;
- une sorte de journaux d'audit ;
- âŠ
Si vous avez une application Java et qu'il n'est pas nĂ©cessaire/possible d'utiliser un cluster Kafka, il existe Ă©galement la possibilitĂ© de travailler via le . L'avantage Ă©vident est que vous pouvez vous passer d'infrastructures supplĂ©mentaires (sous la forme d'un connecteur et de Kafka). Cependant, cette solution est dĂ©clarĂ©e obsolĂšte (deprecated) depuis la version 1.1 et n'est plus recommandĂ©e pour une utilisation (son support pourrait ĂȘtre retirĂ© dans de futures versions).
Dans cet article, nous aborderons l'architecture recommandée par les développeurs, qui garantit la tolérance aux pannes et la capacité d'évoluer.
Configuration du connecteur
Pour commencer Ă suivre les changements de la valeur la plus importante â les donnĂ©es â nous aurons besoin de :
- une source de donnĂ©es, qui peut ĂȘtre MySQL Ă partir de la version 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- un cluster Apache Kafka ;
- une instance Kafka Connect (versions 1.x, 2.x) ;
- un connecteur Debezium configuré.
Les travaux sur les deux premiers points, c'est-à -dire le processus d'installation de la DBMS et d'Apache Kafka, ne sont pas abordés dans cet article. Toutefois, pour ceux qui souhaitent tout déployer dans un bac à sable, le dépÎt officiel contenant des exemples propose un .
Nous allons nous concentrer davantage sur les deux derniers points.
0. Kafka Connect
Ici et dans la suite de l'article, tous les exemples de configuration sont examinés dans le contexte de l'image Docker distribuée par les développeurs de Debezium. Elle contient tous les fichiers de plugins nécessaires (connecteurs) et prévoit la configuration de Kafka Connect via des variables d'environnement.
Si l'utilisation de Kafka Connect de Confluent est envisagĂ©e, il faudra ajouter vous-mĂȘme les plugins des connecteurs nĂ©cessaires dans le rĂ©pertoire spĂ©cifiĂ© dans plugin.path ou dĂ©fini via la variable d'environnement CLASSPATH. Les paramĂštres du worker Kafka Connect et des connecteurs sont dĂ©finis via des fichiers de configuration, qui sont passĂ©s en tant qu'arguments Ă la commande de dĂ©marrage du worker. Pour plus de dĂ©tails, voir .
L'ensemble du processus de configuration de Debezium avec un connecteur se fait en deux étapes. Examinons chacune d'elles :
1. Configuration du framework Kafka Connect
Pour le streaming de donnĂ©es vers un cluster Apache Kafka, des paramĂštres spĂ©cifiques doivent ĂȘtre dĂ©finis dans le framework Kafka Connect, tels que :
- les paramĂštres de connexion au cluster,
- les noms des topics oĂč sera stockĂ©e la configuration du connecteur,
- le nom du groupe dans lequel le connecteur est exécuté (en cas d'utilisation du mode distribué).
L'image Docker officielle du projet prend en charge la configuration via des variables d'environnement â et c'est ce que nous allons utiliser. Donc, tĂ©lĂ©chargeons l'image :
docker pull debezium/connectL'ensemble minimal de variables d'environnement nécessaire au démarrage du connecteur est le suivant :
-
BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092â liste de dĂ©part des serveurs du cluster Kafka pour obtenir la liste complĂšte des membres du cluster ; -
OFFSET_STORAGE_TOPIC=connector-offsetsâ topic pour stocker les positions oĂč se trouve actuellement le connecteur ; -
CONNECT_STATUS_STORAGE_TOPIC=connector-statusâ topic pour stocker l'Ă©tat du connecteur et de ses tĂąches ; -
CONFIG_STORAGE_TOPIC=connector-configâ topic pour stocker les donnĂ©es de configuration du connecteur et de ses tĂąches ; -
GROUP_ID=1â identifiant du groupe de workers sur lesquels la tĂąche du connecteur peut ĂȘtre exĂ©cutĂ©e ; nĂ©cessaire lors de l'utilisation du mode distribuĂ© (distribuĂ©) mode.
Nous lançons le conteneur avec ces variables :
docker run
-e BOOTSTRAP_SERVERS='kafka-1:9092,kafka-2:9092,kafka-3:9092'
-e GROUP_ID=1
-e CONFIG_STORAGE_TOPIC=my_connect_configs
-e OFFSET_STORAGE_TOPIC=my_connect_offsets
-e STATUS_STORAGE_TOPIC=my_connect_statuses debezium/connect:1.2Remarque sur Avro
Par défaut, Debezium écrit des données au format JSON, ce qui est acceptable pour les environnements de test et les petits volumes de données, mais cela peut poser problÚme dans des bases de données à forte charge. Une alternative au convertisseur JSON est la sérialisation des messages à l'aide de en format binaire, ce qui permet de réduire la charge sur le sous-systÚme I/O d'Apache Kafka.
Pour utiliser Avro, il est nécessaire de déployer un (pour stocker les schémas). Les variables pour le convertisseur seront comme suit :
name: CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL
value: http://kafka-registry-01:8081/
name: CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL
value: http://kafka-registry-01:8081/
name: VALUE_CONVERTER
value: io.confluent.connect.avro.AvroConverterLes dĂ©tails sur l'utilisation d'Avro et la configuration du registre dĂ©passent le cadre de cet article â pour plus de clartĂ©, nous utiliserons JSON.
2. Configuration du connecteur lui-mĂȘme
Nous pouvons maintenant passer à la configuration du connecteur proprement dit, qui lira les données à partir de la source.
Prenons comme exemple des connecteurs pour deux SGBD : PostgreSQL et MongoDB, sur lesquels j'ai de l'expĂ©rience et qui prĂ©sentent des diffĂ©rences (mĂȘme petites, mais dans certains cas â significatives !)
La configuration est dĂ©crite dans la notation JSON et est chargĂ©e dans Kafka Connect par une requĂȘte POST.
2.1. PostgreSQL
Exemple de configuration du connecteur pour PostgreSQL :
{
"name": "pg-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"database.hostname": "127.0.0.1",
"database.port": "5432",
"database.user": "debezium",
"database.password": "definitelynotpassword",
"database.dbname": "dbname",
"database.server.name": "pg-dev",
"table.include.list": "public.(.*)",
"heartbeat.interval.ms": "5000",
"slot.name": "dbname_debezium",
"publication.name": "dbname_publication",
"transforms": "AddPrefix",
"transforms.AddPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.AddPrefix.regex": "pg-dev.public.(.*)",
"transforms.AddPrefix.replacement": "data.cdc.dbname"
}
}Le principe de fonctionnement du connecteur aprĂšs cette configuration est plutĂŽt simple :
- Lors du premier dĂ©marrage, il se connecte Ă la base de donnĂ©es indiquĂ©e dans la configuration et fonctionne en mode initial snapshot, envoyant dans Kafka l'ensemble initial des donnĂ©es rĂ©cupĂ©rĂ©es par une requĂȘte
SELECT * FROM table_name. - Une fois l'initialisation terminée, le connecteur passe en mode lecture des changements à partir des fichiers WAL de PostgreSQL.
Concernant les options utilisées :
-
nomâ nom du connecteur utilisĂ© pour la configuration dĂ©crite ci-dessous ; ce nom sera ensuite utilisĂ© pour interagir avec le connecteur (c'est-Ă -dire vĂ©rifier le statut / redĂ©marrer / mettre Ă jour la configuration) via l'API REST de Kafka Connect ; -
connector.classâ classe du connecteur SGBD qui sera utilisĂ©e par le connecteur configurĂ© ; -
plugin.nameâ nom du plugin pour le dĂ©codage logique des donnĂ©es Ă partir des fichiers WAL. Les options disponibles sontwal2json,decoderbuffsetpgoutput. Les deux premiers nĂ©cessitent l'installation des extensions correspondantes dans le SGBD, tandis quepgoutputpour PostgreSQL version 10 et supĂ©rieure, aucune manipulation supplĂ©mentaire n'est nĂ©cessaire ; -
database.*â options pour se connecter Ă la base de donnĂ©es, oĂčdatabase.server.nameâ nom de l'instance PostgreSQL, utilisĂ© pour former le nom du sujet dans le cluster Kafka ; -
table.include.listâ liste des tables dont nous souhaitons suivre les modifications ; elle est spĂ©cifiĂ©e au formatschema.table_name; ne peut pas ĂȘtre utilisĂ©e avectable.exclude.list; -
heartbeat.interval.msâ intervalle (en millisecondes) avec lequel le connecteur envoie des messages heartbeat dans un sujet spĂ©cial ; -
heartbeat.action.queryâ requĂȘte qui sera exĂ©cutĂ©e lors de l'envoi de chaque message heartbeat (option introduite avec la version 1.1) ; -
slot.nameâ nom du slot de rĂ©plication qui sera utilisĂ© par le connecteur ; publication.nameâ nom dans PostgreSQL, utilisĂ©e par le connecteur. Si elle n'existe pas, Debezium essaiera de la crĂ©er. Si l'utilisateur sous lequel la connexion est Ă©tablie n'a pas suffisamment de droits pour cette action â le connecteur se terminera avec une erreur ;-
transformsdéfinit comment modifier le nom du sujet cible :-
transforms.AddPrefix.typeindique que nous allons utiliser des expressions réguliÚres ; -
transforms.AddPrefix.regexâ motif par lequel le nom du sujet cible est redĂ©fini ; -
transforms.AddPrefix.replacementâ ce Ă quoi nous redĂ©finissons effectivement.
-
En savoir plus sur le heartbeat et les transforms
Par dĂ©faut, le connecteur envoie des donnĂ©es vers Kafka Ă chaque transaction validĂ©e, et son LSN (Log Sequence Number) est enregistrĂ© dans un sujet de service offset. Mais que se passe-t-il si le connecteur est configurĂ© pour lire non pas l'ensemble de la base de donnĂ©es, mais seulement certaines de ses tables (oĂč les mises Ă jour des donnĂ©es ne se produisent pas souvent) ?
- Le connecteur lira les fichiers WAL et ne détectera pas de validation des transactions dans les tables qu'il surveille.
- Par conséquent, il ne mettra pas à jour sa position actuelle ni dans le sujet, ni dans le slot de réplication.
- Cela, à son tour, conduira à la "rétention" des fichiers WAL sur le disque et au risque d'épuiser tout l'espace disque.
Et ici, les options viennent Ă la rescousse. heartbeat.interval.ms et heartbeat.action.queryL'utilisation de ces options en tandem permet, Ă chaque fois qu'un message heartbeat est envoyĂ©, d'exĂ©cuter une requĂȘte de modification des donnĂ©es dans une table sĂ©parĂ©e. Cela permet de mettre Ă jour en permanence le LSN sur lequel le connecteur se trouve actuellement (dans le slot de rĂ©plication). Cela permet Ă la SGBD de supprimer les fichiers WAL qui ne sont plus nĂ©cessaires. Pour en savoir plus sur le fonctionnement des options, vous pouvez consulter .
Une autre option, qui mĂ©rite une attention particuliĂšre, est transforms. Bien qu'elle concerne plutĂŽt la commoditĂ© et l'esthĂ©tiqueâŠ
Par dĂ©faut, Debezium crĂ©e des topics en se basant sur la politique de nommage suivante : serverName.schemaName.tableName. Cela peut ne pas toujours ĂȘtre pratique. Les options transforms permettent de dĂ©finir Ă l'aide d'expressions rĂ©guliĂšres la liste des tables dont les Ă©vĂ©nements doivent ĂȘtre routĂ©s vers un topic portant un nom spĂ©cifique.
Dans notre configuration, grĂące Ă Pour conclure la description de la configuration du connecteur PostgreSQL, il convient d'aborder les fonctionnalitĂ©s/limitations suivantes de son fonctionnement : Ainsi, nous chargerons notre configuration dans le connecteur : Nous vĂ©rifions que le chargement s'est effectuĂ© avec succĂšs et que le connecteur est lancĂ© : Excellent : il est configurĂ© et prĂȘt Ă l'emploi. Jouons maintenant le consommateur et connectons-nous Ă Kafka, puis ajoutons et modifions une entrĂ©e dans la table : Dans notre topic, cela se manifestera comme suit : Un JSON trĂšs long avec nos modifications Dans les deux cas, les enregistrements se composent de la clĂ© (PK) de l'enregistrement qui a Ă©tĂ© modifiĂ©, et directement de la nature des changements : quel Ă©tait l'enregistrement avant et quel est le rĂ©sultat aprĂšs. Ce connecteur utilise le mĂ©canisme standard de rĂ©plication de MongoDB, en lisant les informations du oplog du nĆud principal de la base de donnĂ©es. De la mĂȘme maniĂšre que le connecteur pour PgSQL dĂ©crit prĂ©cĂ©demment, ici aussi, lors du premier dĂ©marrage, un instantanĂ© primaire des donnĂ©es est effectuĂ©, aprĂšs quoi le connecteur passe en mode de lecture du oplog. Exemple de configuration : Comme vous pouvez le constater, il n'y a pas de nouvelles options par rapport Ă l'exemple prĂ©cĂ©dent, mais le nombre d'options liĂ©es Ă la connexion Ă la base de donnĂ©es et leurs prĂ©fixes a Ă©tĂ© rĂ©duit. ParamĂštres La question de la tolĂ©rance aux pannes et de la haute disponibilitĂ© est plus cruciale que jamais dans notre Ă©poque â surtout quand nous parlons de donnĂ©es et de transactions, et le suivi des changements de donnĂ©es ne doit pas rester Ă l'Ă©cart de cette question. Examinons ce qui pourrait mal se passer et ce qui arrivera Ă Debezium dans chacun des cas. Il existe trois scĂ©narios de pannes : Cependant, il existe des exceptions. Si le connecteur a Ă©tĂ© dĂ©sactivĂ© pendant une longue pĂ©riode (ou n'a pas pu se connecter Ă l'instance MongoDB), et que l'oplog a Ă©tĂ© tournĂ© pendant ce temps, alors, lors de la rĂ©tablissement de la connexion, le connecteur continuera impassiblement Ă lire les donnĂ©es depuis la premiĂšre position disponible, ce qui signifie qu'une partie des donnĂ©es arrivera dans Kafka. ne . Debezium â c'est ma premiĂšre expĂ©rience avec les systĂšmes CDC et dans l'ensemble, elle est trĂšs positive. Le projet est attirant grĂące Ă la prise en charge des principales bases de donnĂ©es, la simplicitĂ© de configuration, le support de la mise en cluster et une communautĂ© active. Pour ceux qui s'intĂ©ressent Ă la pratique, je recommande de consulter les guides pour et . ComparĂ© au connecteur JDBC pour Kafka Connect, le principal avantage de Debezium est que les modifications sont lues Ă partir des journaux de la SGBD, ce qui permet d'obtenir des donnĂ©es avec un minimum de latence. Le connecteur JDBC (fourni par Kafka Connect) effectue des requĂȘtes sur la table suivie Ă des intervalles fixes et (pour cette mĂȘme raison) ne gĂ©nĂšre pas de messages lors de la suppression de donnĂ©es (comment peut-on demander des donnĂ©es qui n'existent pas ?). Pour des tĂąches similaires, vous pouvez considĂ©rer les solutions suivantes (en plus de Debezium) : Lisez aussi dans notre blog : Source : habr.comtransforms , il se passe ce qui suit : tous les Ă©vĂ©nements CDC de la base de donnĂ©es suivie iront dans un topic nommĂ© data.cdc.dbname. Sinon (sans ces paramĂštres), Debezium crĂ©erait par dĂ©faut un topic pour chaque table sous la forme : pg-dev.public..
Restrictions du connecteur
Application de la configuration
curl -i -X POST -H "Accept:application/json"
-H "Content-Type:application/json" http://localhost:8083/connectors/
-d @pg-con.json$ curl -i http://localhost:8083/connectors/pg-connector/status
HTTP/1.1 200 OK
Date: Thu, 17 Sep 2020 20:19:40 GMT
Content-Type: application/json
Content-Length: 175
Server: Jetty(9.4.20.v20190813)
{"name":"pg-connector","connector":{"state":"RUNNING","worker_id":"172.24.0.5:8083"},"tasks":[{"id":0,"state":"RUNNING","worker_id":"172.24.0.5:8083"}],"type":"source"}$ kafka/bin/kafka-console-consumer.sh
--bootstrap-server kafka:9092
--from-beginning
--property print.key=true
--topic data.cdc.dbname
postgres=# insert into customers (id, first_name, last_name, email) values (1005, 'foo', 'bar', 'foo@bar.com');
INSERT 0 1
postgres=# update customers set first_name = 'egg' where id = 1005;
UPDATE 1{
"schema":{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
}
],
"optional":false,
"name":"data.cdc.dbname.Key"
},
"payload":{
"id":1005
}
}{
"schema":{
"type":"struct",
"fields":[
{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
},
{
"type":"string",
"optional":false,
"field":"first_name"
},
{
"type":"string",
"optional":false,
"field":"last_name"
},
{
"type":"string",
"optional":false,
"field":"email"
}
],
"optional":true,
"name":"data.cdc.dbname.Value",
"field":"before"
},
{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
},
{
"type":"string",
"optional":false,
"field":"first_name"
},
{
"type":"string",
"optional":false,
"field":"last_name"
},
{
"type":"string",
"optional":false,
"field":"email"
}
],
"optional":true,
"name":"data.cdc.dbname.Value",
"field":"after"
},
{
"type":"struct",
"fields":[
{
"type":"string",
"optional":false,
"field":"version"
},
{
"type":"string",
"optional":false,
"field":"connector"
},
{
"type":"string",
"optional":false,
"field":"name"
},
{
"type":"int64",
"optional":false,
"field":"ts_ms"
},
{
"type":"string",
"optional":true,
"name":"io.debezium.data.Enum",
"version":1,
"parameters":{
"allowed":"true,last,false"
},
"default":"false",
"field":"snapshot"
},
{
"type":"string",
"optional":false,
"field":"db"
},
{
"type":"string",
"optional":false,
"field":"schema"
},
{
"type":"string",
"optional":false,
"field":"table"
},
{
"type":"int64",
"optional":true,
"field":"txId"
},
{
"type":"int64",
"optional":true,
"field":"lsn"
},
{
"type":"int64",
"optional":true,
"field":"xmin"
}
],
"optional":false,
"name":"io.debezium.connector.postgresql.Source",
"field":"source"
},
{
"type":"string",
"optional":false,
"field":"op"
},
{
"type":"int64",
"optional":true,
"field":"ts_ms"
},
{
"type":"struct",
"fields":[
{
"type":"string",
"optional":false,
"field":"id"
},
{
"type":"int64",
"optional":false,
"field":"total_order"
},
{
"type":"int64",
"optional":false,
"field":"data_collection_order"
}
],
"optional":true,
"field":"transaction"
}
],
"optional":false,
"name":"data.cdc.dbname.Envelope"
},
"payload":{
"before":null,
"after":{
"id":1005,
"first_name":"foo",
"last_name":"bar",
"email":"foo@bar.com"
},
"source":{
"version":"1.2.3.Final",
"connector":"postgresql",
"name":"dbserver1",
"ts_ms":1600374991648,
"snapshot":"false",
"db":"postgres",
"schema":"public",
"table":"customers",
"txId":602,
"lsn":34088472,
"xmin":null
},
"op":"c",
"ts_ms":1600374991762,
"transaction":null
}
}{
"schema":{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
}
],
"optional":false,
"name":"data.cdc.dbname.Key"
},
"payload":{
"id":1005
}
}{
"schema":{
"type":"struct",
"fields":[
{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
},
{
"type":"string",
"optional":false,
"field":"first_name"
},
{
"type":"string",
"optional":false,
"field":"last_name"
},
{
"type":"string",
"optional":false,
"field":"email"
}
],
"optional":true,
"name":"data.cdc.dbname.Value",
"field":"before"
},
{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
},
{
"type":"string",
"optional":false,
"field":"first_name"
},
{
"type":"string",
"optional":false,
"field":"last_name"
},
{
"type":"string",
"optional":false,
"field":"email"
}
],
"optional":true,
"name":"data.cdc.dbname.Value",
"field":"after"
},
{
"type":"struct",
"fields":[
{
"type":"string",
"optional":false,
"field":"version"
},
{
"type":"string",
"optional":false,
"field":"connector"
},
{
"type":"string",
"optional":false,
"field":"name"
},
{
"type":"int64",
"optional":false,
"field":"ts_ms"
},
{
"type":"string",
"optional":true,
"name":"io.debezium.data.Enum",
"version":1,
"parameters":{
"allowed":"true,last,false"
},
"default":"false",
"field":"snapshot"
},
{
"type":"string",
"optional":false,
"field":"db"
},
{
"type":"string",
"optional":false,
"field":"schema"
},
{
"type":"string",
"optional":false,
"field":"table"
},
{
"type":"int64",
"optional":true,
"field":"txId"
},
{
"type":"int64",
"optional":true,
"field":"lsn"
},
{
"type":"int64",
"optional":true,
"field":"xmin"
}
],
"optional":false,
"name":"io.debezium.connector.postgresql.Source",
"field":"source"
},
{
"type":"string",
"optional":false,
"field":"op"
},
{
"type":"int64",
"optional":true,
"field":"ts_ms"
},
{
"type":"struct",
"fields":[
{
"type":"string",
"optional":false,
"field":"id"
},
{
"type":"int64",
"optional":false,
"field":"total_order"
},
{
"type":"int64",
"optional":false,
"field":"data_collection_order"
}
],
"optional":true,
"field":"transaction"
}
],
"optional":false,
"name":"data.cdc.dbname.Envelope"
},
"payload":{
"before":{
"id":1005,
"first_name":"foo",
"last_name":"bar",
"email":"foo@bar.com"
},
"after":{
"id":1005,
"first_name":"egg",
"last_name":"bar",
"email":"foo@bar.com"
},
"source":{
"version":"1.2.3.Final",
"connector":"postgresql",
"name":"dbserver1",
"ts_ms":1600375609365,
"snapshot":"false",
"db":"postgres",
"schema":"public",
"table":"customers",
"txId":603,
"lsn":34089688,
"xmin":null
},
"op":"u",
"ts_ms":1600375609778,
"transaction":null
}
}
INSERT: la valeur avant (before) est Ă©gale Ă null, et aprĂšs, c'est la chaĂźne qui a Ă©tĂ© insĂ©rĂ©e. UPDATE: dans payload.before l'Ă©tat prĂ©cĂ©dent de la ligne est affichĂ©, et dans payload.after â le nouvel Ă©tat avec la nature des changements.2.2 MongoDB
{
"name": "mp-k8s-mongo-connector",
"config": {
"connector.class": "io.debezium.connector.mongodb.MongoDbConnector",
"tasks.max": "1",
"mongodb.hosts": "MainRepSet/mongo:27017",
"mongodb.name": "mongo",
"mongodb.user": "debezium",
"mongodb.password": "dbname",
"database.whitelist": "db_1,db_2",
"transforms": "AddPrefix",
"transforms.AddPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.AddPrefix.regex": "mongo.([a-zA-Z_0-9]*).([a-zA-Z_0-9]*)",
"transforms.AddPrefix.replacement": "data.cdc.mongo_$1"
}
}transforms cette fois-ci font ce qui suit : transforment le nom du sujet cible selon le schéma <server_name>.<db_name>.<collection_name> dans data.cdc.mongo_<db_name>.Résilience
Conclusion
P.S.
