
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?
è un rappresentante della categoria di software CDC (), ovvero un insieme di connettori per diversi DBMS compatibili con il framework Apache Kafka Connect.
Questo 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:

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 . 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:
- un'origine dati, che può essere MySQL dalla versione 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- cluster Apache Kafka;
- istanza Kafka Connect (versioni 1.x, 2.x);
- 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 .
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 .
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\/connectIl 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.2Nota 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 un formato binario, che contribuisce a ridurre il carico sulla sottosistema I/O in Apache Kafka.
Per utilizzare Avro è necessario distribuire un (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.AvroConverterI 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 sceltawal2json,decoderbuffsepgoutput. I primi due richiedono l'installazione delle rispettive estensioni nel DBMS, mentrepgoutputper PostgreSQL versione 10 e superiori non richiede ulteriori manovre; -
database.*– opzioni per la connessione al DB, dovedatabase.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 formatoschema.table_name; non può essere utilizzato insieme atable.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 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;-
transformsdefinisce come modificare il nome del topic di destinazione:-
transforms.AddPrefix.typeindica 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 .
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 Alla fine della descrizione della configurazione del connettore per PostgreSQL, vale la pena menzionare le seguenti caratteristiche/limitazioni del suo funzionamento: Quindi, carichiamo la nostra configurazione nel connettore: Controlliamo che il caricamento sia andato a buon fine e che il connettore sia avviato: 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: Nel nostro argomento questo si rifletterà nel seguente modo: JSON molto lungo con le nostre modifiche 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. 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: 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 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: 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 . 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 e . 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): Leggi anche nel nostro blog: Fonte: habr.comtransforms 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
Applicazione della configurazione
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
}
}
INSERISCI: valore prima (before) uguale null, e dopo — la stringa che è stata inserita. UPDATE: in payload.before visualizza lo stato precedente della riga, mentre in payload.after — è il nuovo con la sostanza delle modifiche.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 questa volta fanno quanto segue: trasformano il nome del topic di destinazione dal modello .. in data.cdc.mongo_.Resilienza
Conclusione
P.S.
