Introduzione a Debezium — CDC per Apache Kafka

Introduzione a Debezium — CDC per Apache Kafka

Nel mio lavoro mi imbatto spesso in nuove soluzioni tecniche/prodotti software, di cui ci sono poche informazioni su internet in lingua russa. Con questo articolo cercherò di colmare una di queste lacune con un esempio dalla mia recente pratica, quando è stato necessario configurare l'invio di eventi CDC da due popolari DBMS (PostgreSQL e MongoDB) a un cluster Kafka utilizzando Debezium. Spero che questo articolo riassuntivo, apparso come risultato del lavoro svolto, possa essere utile anche ad altri.

Cos'è Debezium e che cos'è il CDC?

Debezium è un rappresentante della categoria di software CDC (Capture Data Change), ovvero un insieme di connettori per diversi DBMS compatibili con il framework Apache Kafka Connect.

Questo Progetto open source, che utilizza la licenza Apache License v2.0 ed è sponsorizzato da Red Hat. Lo sviluppo è iniziato nel 2016 e attualmente supporta ufficialmente i seguenti DBMS: MySQL, PostgreSQL, MongoDB, SQL Server. Esistono anche connettori per Cassandra e Oracle, ma attualmente sono in stato di "accesso anticipato" e le nuove versioni non garantiscono la retrocompatibilità.

Se confrontiamo il CDC con l'approccio tradizionale (quando l'applicazione legge direttamente i dati dal DBMS), i suoi principali vantaggi includono la realizzazione dello streaming delle modifiche dei dati a livello di righe con bassa latenza, alta affidabilità e disponibilità. Gli ultimi due punti vengono raggiunti grazie all'utilizzo di un cluster Kafka come archivio per gli eventi CDC.

Inoltre, tra i vantaggi si può includere il fatto che per l'archiviazione degli eventi viene utilizzato un modello unificato, pertanto l'applicazione finale non dovrà preoccuparsi delle complessità dell'operazione di diversi DBMS.

Infine, grazie all'utilizzo di un broker di messaggi si apre la possibilità di scalare orizzontalmente le applicazioni che monitorano le modifiche ai dati. In questo modo, l'impatto sulla fonte dei dati è minimizzato, poiché i dati non vengono acquisiti direttamente dal DBMS, ma dal cluster Kafka.

Architettura di Debezium

L'uso di Debezium si riduce a uno schema molto semplice:

DBMS (come fonte dei dati) → connettore in Kafka Connect → Apache Kafka → consumer

Come illustrazione, cito uno schema dal sito del progetto:

Introduzione a Debezium — CDC per Apache Kafka

Tuttavia, questo schema non mi piace molto, poiché dà l'impressione che sia possibile solo utilizzare un connettore sink.

In realtà, la situazione è diversa: il riempimento del vostro Data Lake (l'ultimo anello nello schema sopra) — non è l'unico modo per utilizzare Debezium. Gli eventi inviati in Apache Kafka possono essere utilizzati dalle vostre applicazioni per risolvere diverse situazioni. Ad esempio:

  • rimozione di dati obsoleti dalla cache;
  • invio di notifiche;
  • aggiornamenti degli indici di ricerca;
  • una sorta di registri di audit;

Nel caso in cui abbiate un'applicazione Java e non sia necessario/possibile utilizzare un cluster Kafka, c'è anche la possibilità di lavorare tramite embedded-connector. Un chiaro vantaggio è che con esso si può rinunciare a infrastrutture aggiuntive (sotto forma di connettore e Kafka). Tuttavia, questa soluzione è stata dichiarata obsoleta (deprecated) dalla versione 1.1 e non è più raccomandata per l'uso (nelle future release il supporto potrebbe essere rimosso).

In questo articolo verrà discussa l'architettura consigliata dagli sviluppatori, che garantisce resilienza e scalabilità.

Configurazione del connettore

Per iniziare a monitorare le modifiche al valore più importante — i dati — avremo bisogno di:

  1. un'origine dati, che può essere MySQL dalla versione 5.7, PostgreSQL 9.6+, MongoDB 3.2+ (l'elenco completo);
  2. cluster Apache Kafka;
  3. istanza Kafka Connect (versioni 1.x, 2.x);
  4. un connettore Debezium configurato.

Le operazioni sui primi due punti, ovvero il processo di installazione del DBMS e di Apache Kafka, esulano dallo scopo di questo articolo. Tuttavia, per chi desidera implementare tutto in un ambiente sandbox, nel repository ufficiale degli esempi è disponibile un docker-compose.yaml.

Noi ci concentreremo più dettagliatamente sugli ultimi due punti.

0. Kafka Connect

Qui e in seguito nell'articolo, tutti gli esempi di configurazione saranno considerati nel contesto dell'immagine Docker distribuita dagli sviluppatori di Debezium. Essa contiene tutti i file plugin necessari (connettori) e prevede la configurazione di Kafka Connect tramite variabili d'ambiente.

Nel caso si preveda di utilizzare Kafka Connect da Confluent, sarà necessario aggiungere manualmente i plugin dei connettori necessari nella directory specificata in plugin.path o definita tramite la variabile d'ambiente CLASSPATH. Le impostazioni del worker di Kafka Connect e dei connettori vengono definite tramite file di configurazione, che vengono passati come argomenti al comando di avvio del worker. Maggiori dettagli sono disponibili in documentazione.

L'intero processo di configurazione di Debezium con il connettore avviene in due fasi. Esaminiamo ciascuna di esse:

1. Configurazione del framework Kafka Connect

Per lo streaming dei dati nel cluster Apache Kafka nel framework Kafka Connect vengono impostati parametri specifici, come ad esempio:

  • parametri di connessione al cluster,
  • nomi dei topic in cui verrà memorizzata direttamente la configurazione del connettore,
  • nome del gruppo in cui è in esecuzione il connettore (in caso di utilizzo della modalità distribuita).

L'immagine ufficiale Docker del progetto supporta la configurazione tramite variabili d'ambiente — ecco come procederemo. Quindi, scarichiamo l'immagine:

docker pull debezium\/connect

Il set minimo di variabili d'ambiente necessarie per avviare il connettore è il seguente:

  • BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092 — elenco iniziale dei server del cluster Kafka per ottenere l'elenco completo dei membri del cluster;
  • OFFSET_STORAGE_TOPIC=connector-offsets — topic utilizzato per memorizzare le posizioni in cui si trova attualmente il connettore;
  • CONNECT_STATUS_STORAGE_TOPIC=connector-status — topic per memorizzare lo stato del connettore e delle sue attività;
  • CONFIG_STORAGE_TOPIC=connector-config — topic per memorizzare i dati di configurazione del connettore e delle sue attività;
  • GROUP_ID=1 — identificatore del gruppo di lavoratori su cui può essere eseguita l'attività del connettore; necessario per l'uso della modalità distribuita (distribuita) modalità.

Avviamo il container con queste variabili:

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

Nota su Avro

Per impostazione predefinita, Debezium scrive i dati in formato JSON, che è accettabile per sandbox e piccoli volumi di dati, ma può diventare problematico in database ad alta richiesta. Un'alternativa al convertitore JSON è la serializzazione dei messaggi tramite Avro un formato binario, che contribuisce a ridurre il carico sulla sottosistema I/O in Apache Kafka.

Per utilizzare Avro è necessario distribuire un schema-registry (per la memorizzazione degli schemi). Le variabili per il convertitore appariranno come segue:

name: CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL
value: http:\/\/kafka-registry-01:8081\/\nname: CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL
value: http:\/\/kafka-registry-01:8081\/\nname: VALUE_CONVERTER   
value: io.confluent.connect.avro.AvroConverter

I dettagli sull'uso di Avro e sulla configurazione del registry vanno oltre l'argomento di questo articolo — per chiarezza, utilizzeremo JSON.

2. Configurazione del connettore

Ora possiamo passare direttamente alla configurazione del connettore stesso, che leggerà i dati dalla sorgente.

Prendiamo ad esempio i connettori per due DBMS: PostgreSQL e MongoDB, su cui ho esperienza e che presentano alcune differenze (anche se piccole, ma in alcuni casi – significative!).

La configurazione è descritta in notazione JSON e viene caricata in Kafka Connect tramite una richiesta POST.

2.1. PostgreSQL

Esempio di configurazione del connettore per 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"
  }
}

Il principio di funzionamento del connettore dopo tale configurazione è piuttosto semplice:

  • Alla prima esecuzione si connette al database specificato nella configurazione e viene avviato in modalità initial snapshot, inviando a Kafka un set iniziale di dati ottenuti mediante il comando SELECT * FROM table_name.
  • Una volta completata l'inizializzazione, il connettore passa alla modalità di lettura delle modifiche dai file WAL di PostgreSQL.

Sulle opzioni utilizzate:

  • name – nome del connettore per il quale viene utilizzata la configurazione descritta di seguito; in seguito, questo nome viene usato per interagire con il connettore (ovvero, per controllare lo stato/riavviare/aggiornare la configurazione) tramite l'API REST di Kafka Connect;
  • connector.class – classe del connettore DBMS che sarà utilizzato dal connettore configurabile;
  • plugin.name – nome del plugin per la decodifica logica dei dati dai file WAL. Sono disponibili per la scelta wal2json, decoderbuffs e pgoutput. I primi due richiedono l'installazione delle rispettive estensioni nel DBMS, mentre pgoutput per PostgreSQL versione 10 e superiori non richiede ulteriori manovre;
  • database.* – opzioni per la connessione al DB, dove database.server.name – nome dell'istanza PostgreSQL utilizzato per formare il nome del topic nel cluster Kafka;
  • table.include.list – lista delle tabelle in cui desideriamo monitorare le modifiche; viene specificata nel formato schema.table_name; non può essere utilizzato insieme a table.exclude.list;
  • heartbeat.interval.ms — intervallo (in millisecondi) con cui il connettore invia messaggi heartbeat a un topic speciale;
  • heartbeat.action.query — query che verrà eseguita al momento dell'invio di ogni messaggio heartbeat (opzione disponibile dalla versione 1.1);
  • slot.name — nome dello slot di replica che verrà utilizzato dal connettore;
  • publication.name — nome pubblicazione in PostgreSQL, utilizzato dal connettore. Se non esiste, Debezium tenterà di crearlo. Se l'utente con cui ci si connette non ha i permessi sufficienti per questa operazione, il connettore terminerà con un errore;
  • transforms definisce come modificare il nome del topic di destinazione:
    • transforms.AddPrefix.type indica che verranno utilizzate espressioni regolari;
    • transforms.AddPrefix.regex — maschera mediante la quale viene ridefinito il nome del topic di destinazione;
    • transforms.AddPrefix.replacement — ciò a cui stiamo effettuando la ridefinizione.

Ulteriori informazioni su heartbeat e transforms

Di default, il connettore invia dati a Kafka per ogni transazione confermata, registrando il suo LSN (Log Sequence Number) in un topic di servizio offset. Ma cosa succede se il connettore è configurato per leggere non l'intero database, ma solo alcune delle sue tabelle (dove l'aggiornamento dei dati non avviene frequentemente)?

  • Il connettore leggerà i file WAL e non troverà in essi le transazioni confermate per le tabelle che sta monitorando.
  • Pertanto, non aggiornerà la sua posizione attuale né nel topic né nello slot di replica.
  • Questo, a sua volta, porterà a un 'blocco' dei file WAL sul disco e al possibile esaurimento dello spazio su disco.

E qui entrano in gioco le opzioni heartbeat.interval.ms e heartbeat.action.query. L'uso di queste opzioni insieme consente di eseguire una query per modificare i dati in una tabella separata ogni volta che si invia un messaggio heartbeat. Ciò aggiorna costantemente l’LSN su cui si trova attualmente il connettore (nello slot di replica). Questo permette al DBMS di eliminare i file WAL non più necessari. Maggiori dettagli sul funzionamento delle opzioni possono essere trovati in documentazione.

Un'altra opzione meritevole di maggiore attenzione è transforms. Anche se riguarda più il comfort e l’estetica…

Di default, Debezium crea topic seguendo la seguente politica di denominazione: serverName.schemaName.tableName. Questo non sempre può essere conveniente. Opzioni transforms È possibile utilizzare espressioni regolari per determinare un elenco di tabelle, gli eventi dai quali devono essere indirizzati in un argomento con un nome specifico.

Nella nostra configurazione grazie a transforms si verifica quanto segue: tutti gli eventi CDC dal database monitorato verranno inviati all'argomento con il nome data.cdc.dbname. Altrimenti (senza queste impostazioni) Debezium creerebbe per impostazione predefinita un argomento per ogni tabella di tipo: pg-dev.public.

.

Limitazioni del connettore

Alla fine della descrizione della configurazione del connettore per PostgreSQL, vale la pena menzionare le seguenti caratteristiche/limitazioni del suo funzionamento:

  1. La funzionalità del connettore per PostgreSQL si basa sul concetto di decodifica logica. Pertanto, non tiene traccia delle query che modificano la struttura del database (DDL) — conseguentemente, queste informazioni non saranno presenti negli argomenti.
  2. Poiché vengono utilizzati slot di replica, la connessione del connettore è possibile solo a un'istanza primaria del DBMS.
  3. Se all'utente con cui il connettore si connette al database sono stati concessi solo diritti di lettura, prima del primo avvio sarà necessario creare manualmente uno slot di replica e una pubblicazione nel database.

Applicazione della configurazione

Quindi, carichiamo la nostra configurazione nel connettore:

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

Controlliamo che il caricamento sia andato a buon fine e che il connettore sia avviato:

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

Ottimo: è configurato e pronto per l'uso. Ora fingiamo di essere un consumer e connettiamoci a Kafka, dopodiché aggiungiamo e modifichiamo una registrazione nella tabella:

$ 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

Nel nostro argomento questo si rifletterà nel seguente modo:

JSON molto lungo con le nostre modifiche

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

In entrambi i casi, le registrazioni consistono nella chiave (PK) della registrazione che è stata modificata e nella sostanza stessa delle modifiche: com'era la registrazione prima e com'era dopo.

  • Nel caso di INSERISCI: valore prima (before) uguale null, e dopo — la stringa che è stata inserita.
  • Nel caso di UPDATE: in payload.before visualizza lo stato precedente della riga, mentre in payload.after — è il nuovo con la sostanza delle modifiche.

2.2 MongoDB

Questo connettore utilizza il meccanismo di replica standard di MongoDB, leggendo le informazioni dal oplog del nodo primary del DBMS.

Analogamente al connettore per PgSQL già descritto, anche qui, al primo avvio, viene scattato uno snapshot primario dei dati, dopodiché il connettore passa alla modalità di lettura dell'oplog.

Esempio di configurazione:

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

Come si può notare, non ci sono nuove opzioni rispetto all'esempio precedente, ma è diminuito solo il numero delle opzioni che riguardano la connessione al DB e i loro prefissi.

Impostazioni transforms questa volta fanno quanto segue: trasformano il nome del topic di destinazione dal modello .. in data.cdc.mongo_.

Resilienza

La questione della tolleranza ai guasti e dell'alta disponibilità è più attuale che mai — specialmente quando parliamo di dati e transazioni, e il monitoraggio delle modifiche ai dati non rimane estraneo a questa questione. Vediamo cosa può fondamentalmente andare storto e cosa accadrà a Debezium in ognuno dei casi.

Ci sono tre opzioni di guasto:

  1. Guasto di Kafka Connect. Se Connect è configurato per funzionare in modalità distribuita, è necessario che più worker abbiano lo stesso group.id. In tal modo, in caso di guasto di uno di essi, il connettore verrà riavviato su un altro worker e continuerà a leggere dalla posizione dell'ultima commit nel topic in Kafka.
  2. Perdita di connettività con il cluster Kafka. Il connettore semplicemente interromperà la lettura dalla posizione che non è riuscito a inviare a Kafka e cercherà periodicamente di reinviarla fino a quando il tentativo non avrà successo.
  3. Inaccessibilità della fonte di dati. Il connettore tenterà di riconnettersi alla sorgente secondo la configurazione. Per impostazione predefinita, ci saranno 16 tentativi utilizzando il backoff esponenziale. Dopo il sedicesimo tentativo fallito, il task sarà contrassegnato come fallito e richiederà un riavvio manuale tramite l'interfaccia REST di Kafka Connect.
    • Nel caso di PostgreSQL I dati non andranno persi, poiché l'uso delle slot di replica impedirà l'eliminazione dei file WAL non letti dal connettore. In questo caso, c'è anche un rovescio della medaglia: se per un periodo prolungato la connessione di rete tra il connettore e il DBMS viene interrotta, c'è il rischio che lo spazio su disco si esaurisca, il che potrebbe portare a un completo fallimento del DBMS.
    • Nel caso di MySQL I file binlog possono essere ruotati dal DBMS stesso prima che la connettività venga ripristinata. Ciò farà sì che il connettore entri in stato di fallito e sarà necessario un riavvio nel formato di snapshot iniziale per continuare a leggere dai binlog.
    • Su MongoDB. La documentazione afferma: il comportamento del connettore nel caso in cui i file di registro/oplog siano stati eliminati e il connettore non possa continuare a leggere dalla posizione in cui si era fermato è lo stesso per tutti i DBMS. In sintesi, il connettore entrerà in stato di fallito e richiederà un riavvio in modalità initial snapshot.

      Tuttavia, ci sono eccezioni. Se il connettore è rimasto disconnesso per un lungo periodo di tempo (o non è riuscito a contattare l'istanza di MongoDB), e l'oplog è stato ruotato nel frattempo, quando la connessione viene ripristinata, il connettore continuerà serenamente a leggere i dati dalla prima posizione disponibile, il che farà sì che parte dei dati vada a finire in Kafka non .

Conclusione

Debezium è la mia prima esperienza con sistemi CDC e in generale è stata molto positiva. Il progetto conquista per il supporto ai principali DBMS, la semplicità di configurazione, il supporto per la clusterizzazione e la comunità attiva. A chi è interessato alla pratica consiglio di consultare le guide per Kafka Connect e Debezium.

Rispetto al connettore JDBC per Kafka Connect, il principale vantaggio di Debezium è che le modifiche vengono lette dai registri del DBMS, consentendo di ricevere i dati con la minima latenza. Il connettore JDBC (nella fornitura di Kafka Connect) esegue richieste sulla tabella monitorata a intervalli fissi e (per lo stesso motivo) non genera messaggi in caso di eliminazione dei dati (come può essere richiesto dati che non esistono?).

Per affrontare compiti simili, si possono considerare le seguenti soluzioni (oltre a Debezium):

P.S.

Leggi anche nel nostro blog:

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster