DĂ©couverte de Debezium — CDC pour Apache Kafka

DĂ©couverte de Debezium — CDC pour Apache Kafka

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 ?

Debezium est un représentant de la catégorie des logiciels CDC (Capture Data Change), 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 un projet Open Source, 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 :

DĂ©couverte de Debezium — CDC pour Apache Kafka

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 connecteur embedded. 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 :

  1. une source de donnĂ©es, qui peut ĂȘtre MySQL Ă  partir de la version 5.7, PostgreSQL 9.6+, MongoDB 3.2+ (liste complĂšte);
  2. un cluster Apache Kafka ;
  3. une instance Kafka Connect (versions 1.x, 2.x) ;
  4. 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 docker-compose.yaml.

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 documentation.

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/connect

L'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.2

Remarque 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 Avro 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 schema-registry (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.AvroConverter

Les 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 sont wal2json, decoderbuffs et pgoutput. Les deux premiers nĂ©cessitent l'installation des extensions correspondantes dans le SGBD, tandis que pgoutput pour 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 format schema.table_name; ne peut pas ĂȘtre utilisĂ©e avec table.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 de la publication 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 ;
  • transforms dĂ©finit comment modifier le nom du sujet cible :
    • transforms.AddPrefix.type indique 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 documentation.

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 à transforms , 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

Pour conclure la description de la configuration du connecteur PostgreSQL, il convient d'aborder les fonctionnalités/limitations suivantes de son fonctionnement :

  1. La fonctionnalitĂ© du connecteur PostgreSQL repose sur le concept de dĂ©codage logique. Par consĂ©quent, il ne suit pas les requĂȘtes de modification de la structure de la base de donnĂ©es (DDL) — par consĂ©quent, ces donnĂ©es ne seront pas prĂ©sentes dans les topics.
  2. Étant donnĂ© que des slots de rĂ©plication sont utilisĂ©s, la connexion du connecteur est possible uniquement au module principal du SGBD.
  3. Si l'utilisateur sous lequel le connecteur se connecte à la base de données a uniquement des droits en lecture, il sera nécessaire de créer manuellement un slot de réplication et une publication dans la base de données avant le premier lancement.

Application de la configuration

Ainsi, nous chargerons notre configuration dans le connecteur :

curl -i -X POST -H "Accept:application/json" 
  -H  "Content-Type:application/json"  http://localhost:8083/connectors/ 
  -d @pg-con.json

Nous vérifions que le chargement s'est effectué avec succÚs et que le connecteur est lancé :

$ 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"}

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 :

$ 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

Dans notre topic, cela se manifestera comme suit :

Un JSON trĂšs long avec nos modifications

{
  "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
  }
}

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.

  • Dans le cas de INSERT: la valeur avant (before) est Ă©gale Ă  null, et aprĂšs, c'est la chaĂźne qui a Ă©tĂ© insĂ©rĂ©e.
  • Dans le cas de 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

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 :

{
  "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"
        }
  }

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 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

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 :

  1. Panne de Kafka Connect. Si Connect est configurĂ© pour fonctionner en mode distribuĂ©, plusieurs travailleurs doivent avoir le mĂȘme group.id. Ainsi, en cas de panne de l'un d'eux, le connecteur sera redĂ©marrĂ© sur un autre travailleur et continuera la lecture Ă  partir de la derniĂšre position validĂ©e dans le sujet Kafka.
  2. Perte de connectivitĂ© avec le cluster Kafka. Le connecteur va simplement arrĂȘter la lecture Ă  la position qui n'a pas pu ĂȘtre envoyĂ©e Ă  Kafka, et essaiera pĂ©riodiquement de la renvoyer jusqu'Ă  ce que la tentative soit rĂ©ussie.
  3. Inaccessibilité de la source de données. Le connecteur essaiera de se reconnecter à la source selon la configuration. Par défaut, cela représente 16 tentatives utilisant un retour exponentiel. AprÚs la 16e tentative infructueuse, la tùche sera marquée comme échec et nécessitera un redémarrage manuel via l'API REST de Kafka Connect.
    • Dans le cas de PostgreSQL Les donnĂ©es ne seront pas perdues, car l'utilisation de slots de rĂ©plication empĂȘchera la suppression des fichiers WAL non lus par le connecteur. Cependant, il y a un revers Ă  la mĂ©daille : si la connectivitĂ© rĂ©seau entre le connecteur et la base de donnĂ©es est interrompue pendant une longue pĂ©riode, il y a un risque que l'espace disque soit Ă©puisĂ©, ce qui pourrait entraĂźner un Ă©chec total de la base de donnĂ©es.
    • Dans le cas de MySQL Les fichiers binlogs peuvent ĂȘtre tournĂ©s par la base de donnĂ©es avant que la connectivitĂ© soit rĂ©tablie. Cela entraĂźnera un passage du connecteur Ă  un Ă©tat d'Ă©chec, et un redĂ©marrage en mode de snapshot initial sera nĂ©cessaire pour continuer Ă  lire les binlogs.
    • À propos de MongoDB. La documentation stipule que le comportement du connecteur dans le cas oĂč les fichiers de journaux/oplog ont Ă©tĂ© supprimĂ©s et que le connecteur ne peut pas continuer Ă  lire Ă  partir de la position oĂč il s'est arrĂȘtĂ© est identique pour toutes les bases de donnĂ©es. Il s'agit d'un passage du connecteur Ă  un Ă©tat Ă©chec et nĂ©cessitera un redĂ©marrage en mode initial snapshot.

      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 .

Conclusion

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 Kafka Connect et Debezium.

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) :

P.S.

Lisez aussi dans notre blog :

Source : habr.com

Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS đŸ”„ Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster