
Nella mia attività, mi trovo spesso di fronte a nuove soluzioni tecniche/prodotti software di cui ci sono poche informazioni nel web di lingua russa. Con questo articolo cercherò di colmare una di queste lacune con un esempio tratto dalla mia esperienza recente, 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 di panoramica, risultante dal lavoro svolto, sia utile anche ad altri.
Che cos'è Debezium e cos'è in generale il CDC?
è un rappresentante della categoria del software CDC (), e più precisamente è un insieme di connettori per vari 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 e SQL Server. Sono disponibili anche connettori per Cassandra e Oracle, ma attualmente sono in fase di "accesso anticipato", e le nuove versioni non garantiscono la compatibilità retroattiva.
Confrontando il CDC con l'approccio tradizionale (dove l'applicazione legge i dati direttamente dal database), i suoi principali vantaggi includono l'implementazione dello streaming delle modifiche ai dati a livello di riga con bassa latenza, alta affidabilità e disponibilità. Gli ultimi due elementi sono raggiunti grazie all'uso di un cluster Kafka come archivio degli eventi CDC.
Un ulteriore vantaggio è che per la memorizzazione degli eventi viene utilizzato un modello unificato, quindi l'applicazione finale non deve preoccuparsi delle complessità nella gestione di diversi database.
Infine, grazie all'uso di un broker di messaggi, si apre la possibilità di scalare orizzontalmente le applicazioni che monitorano le modifiche nei dati. Ciò limita al minimo l'impatto sulla fonte di dati, poiché i dati non vengono ottenuti direttamente dal database, ma dal cluster Kafka.
Architettura di Debezium
L'uso di Debezium si riduce a uno schema semplice:
Database (come sorgente dei dati) → connettore in Kafka Connect → Apache Kafka → consumer
Come illustrazione, ecco uno schema dal sito del progetto:

Tuttavia, questo schema non mi convince molto, poiché sembra che si possa utilizzare solo 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 ad 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;
- …
Se avete un'applicazione Java e non avete la necessità/la possibilità di utilizzare un cluster Kafka, è anche possibile lavorare tramite . Un vantaggio evidente è che si può rinunciare a un'infrastruttura aggiuntiva (sotto forma di connettore e Kafka). Tuttavia, questa soluzione è stata dichiarata obsoleta (deprecated) dalla versione 1.1 e non è più raccomandata. (Nei futuri rilasci potrebbe essere rimossa la sua supporto).
In questo articolo verrà discussa l'architettura raccomandata dagli sviluppatori, che garantisce resilienza e scalabilità.
Configurazione del connettore
Per iniziare a monitorare le variazioni della nostra risorsa più importante — i dati — avremo bisogno di:
- una fonte di dati, che può essere MySQL dalla versione 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- un cluster Apache Kafka;
- un'istanza Kafka Connect (versioni 1.x, 2.x);
- un connettore Debezium configurato.
I primi due punti, ossia l'installazione del DBMS e di Apache Kafka, vanno oltre l'ambito di questo articolo. Tuttavia, per chi desidera allestire il tutto in un ambiente sandbox, nel repository ufficiale degli esempi c'è un facile .
Ci concentreremo più dettagliatamente sugli ultimi due punti.
0. Kafka Connect
Qui e nelle sezioni successive dell'articolo, tutti gli esempi di configurazione sono presentati nel contesto dell'immagine Docker distribuita dagli sviluppatori di Debezium. Essa contiene tutti i file necessari dei plugin (connettori) e prevede la configurazione di Kafka Connect tramite variabili d'ambiente.
Se si prevede di utilizzare Kafka Connect da Confluent, sarà necessario aggiungere manualmente i plugin dei connettori necessari nella directory indicata in plugin.path o specificata tramite la variabile d'ambiente CLASSPATH. Le impostazioni del worker di Kafka Connect e dei connettori sono 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 si svolge in due fasi. Analizziamo 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:
- parametri di connessione al cluster,
- nomi dei topic in cui sarà memorizzata la configurazione del connettore stesso,
- nome del gruppo in cui è in esecuzione il connettore (in caso di utilizzo della modalità distribuita).
L'immagine Docker ufficiale del progetto supporta la configurazione tramite variabili d'ambiente — e questo è ciò che utilizzeremo. Quindi, scarichiamo l'immagine:
docker pull debezium/connectL'insieme minimo di variabili d'ambiente necessarie per avviare il connettore è il seguente:
-
BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092— un elenco iniziale dei server del cluster Kafka per ottenere l'elenco completo dei membri del cluster; -
OFFSET_STORAGE_TOPIC=connector-offsets— il topic per memorizzare le posizioni in cui si trova attualmente il connettore; -
CONNECT_STATUS_STORAGE_TOPIC=connector-status— argomento per la memorizzazione dello stato del connettore e delle sue attività; -
CONFIG_STORAGE_TOPIC=connector-config— argomento per la memorizzazione dei 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 quando si utilizza la modalità distribuita (distributed) mode.
Avviamo il contenitore 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 ambienti di test e piccole quantità di dati, ma può diventare problematico in database ad alta intensità di carico. Un'alternativa al convertitore JSON è la serializzazione dei messaggi nel formato binario, il che aiuta a ridurre il carico sul 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/
name: CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL
value: http://kafka-registry-01:8081/
name: VALUE_CONVERTER
value: io.confluent.connect.avro.AvroConverterLe informazioni sull'uso di Avro e sulla sua configurazione esulano dallo scopo di questo articolo — per chiarezza, utilizzeremo JSON.
2. Configurazione del connettore
Possiamo ora passare direttamente alla configurazione del connettore, che leggerà i dati dalla sorgente.
Prendiamo come esempio connettori per due DBMS: PostgreSQL e MongoDB, per i quali ho esperienza e per i quali esistono differenze (anche se piccole, in alcuni casi possono essere 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 funzionamento del connettore dopo questa configurazione è piuttosto semplice:
- All'avvio, si collega al database specificato nella configurazione e si avvia in modalità istantanea iniziale, inviando a Kafka un set iniziale di dati ottenuti tramite il comando
SELECT * FROM table_name. - Dopo che l'inizializzazione è completata, il connettore passa alla modalità di lettura delle modifiche dai file WAL di PostgreSQL.
Riguardo le opzioni utilizzate:
-
name— nome del connettore per il quale viene utilizzata la configurazione descritta di seguito; in seguito, questo nome sarà utilizzato per interagire con il connettore (es. controllare lo stato/ripristinare/aggiornare la configurazione) attraverso le API REST di Kafka Connect; -
connector.class— classe di connettore DBMS che sarà utilizzata dal connettore configurabile; -
plugin.name— nome del plugin per la decodifica logica dei dati dai file WAL. Sono disponibili le seguenti opzioni:wal2json,decoderbuffsepgoutput. I primi due richiedono l'installazione delle corrispondenti estensioni nel DBMS, mentrepgoutputper PostgreSQL versione 10 e superiori non richiede ulteriori manipolazioni; -
database.*— opzioni per la connessione al DB, dovedatabase.server.name— nome dell'istanza PostgreSQL, utilizzato per generare il nome del topic nel cluster Kafka; -
table.include.list— elenco delle tabelle in cui vogliamo monitorare le modifiche; specificato nel formatoschema.table_name; non può essere usato insieme atable.exclude.list; -
heartbeat.interval.ms— intervallo (in millisecondi) in cui il connettore invia messaggi heartbeat a un topic speciale; -
heartbeat.action.query— query che verrà eseguita al inviare ogni messaggio heartbeat (questa opzione è stata introdotta dalla versione 1.1); -
slot.name— nome dello slot di replica che sarà utilizzato dal connettore; publication.name— nome in PostgreSQL, usato dal connettore. Se non esiste, Debezium tenterà di crearla. Se l'utente con cui viene effettuata la connessione non ha sufficienti diritti per eseguire questa azione, il connettore terminerà con un errore;-
trasformazionidefinisce come modificare il nome del topic di destinazione:-
transforms.AddPrefix.typeindica che utilizzeremo espressioni regolari; -
transforms.AddPrefix.regex— maschera utilizzata per sovrascrivere il nome del topic di destinazione; -
transforms.AddPrefix.replacement— valore effettivo su cui stiamo effettuando la sovrascrittura.
-
Ulteriori informazioni su heartbeat e trasformazioni
Per impostazione predefinita, il connettore invia i dati a Kafka per ogni transazione confermata, e registra il suo LSN (Log Sequence Number) in un topic di servizio offset. Ma cosa succede se il connettore è configurato per leggere solo una parte delle sue tabelle (in cui l'aggiornamento dei dati non avviene spesso)?
- Il connettore leggerà i file WAL e non scoprirebbe la conferma delle transazioni in quelle 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 "trattenere" i file WAL su disco e a una probabile esaurimento dello spazio su disco.
E qui entrano in gioco le opzioni heartbeat.interval.ms e heartbeat.action.query. L'uso di queste opzioni in coppia consente di eseguire una richiesta di modifica dei dati in una tabella separata ogni volta che viene inviato un messaggio heartbeat. In questo modo, il LSN su cui si trova attualmente il connettore (nello slot di replicazione) viene costantemente aggiornato. Ciò consente al DBMS di eliminare i file WAL che non sono più necessari. Per saperne di più sul funzionamento delle opzioni, puoi consultare .
Un'altra opzione meritevole di maggiore attenzione è trasformazioni. Sebbene riguardi più il comfort e l'estetica...
Per impostazione predefinita, Debezium crea argomenti seguendo la seguente politica di denominazione: serverName.schemaName.tableName. Questo potrebbe non essere sempre conveniente. Con le opzioni trasformazioni è possibile definire, utilizzando espressioni regolari, l'elenco delle tabelle gli eventi delle quali devono essere instradati in un argomento con un nome specifico.
Nella nostra configurazione, grazie a Alla fine della descrizione della configurazione del connettore per PostgreSQL, è opportuno parlare delle seguenti caratteristiche/limitazioni del suo funzionamento: Pertanto, carichiamo la nostra configurazione nel connettore: Verifichiamo che il caricamento sia avvenuto con successo e che il connettore sia partito: Ottimo: è configurato e pronto per l'uso. Ora faremo finta di essere un consumer e ci collegheremo a Kafka, dopo di che aggiungeremo e modificheremo una voce nella tabella: Nel nostro topic questo apparirà come segue: JSON molto lungo con le nostre modifiche In entrambi i casi, le registrazioni consistono in una chiave (PK) della registrazione che è stata modificata e nel merito stesso delle modifiche: quale fosse la registrazione prima e quale è diventata dopo. Questo connettore utilizza il meccanismo standard di replica di MongoDB, leggendo le informazioni dall'oplog del nodo primario del DBMS. Analogamente al connettore descritto per PgSQL, qui viene comunque scattato un'istantanea primaria dei dati al primo avvio, 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 relative alla connessione al database e i loro prefissi. Impostazioni La questione della tolleranza ai guasti e dell'alta disponibilità è più attuale che mai, soprattutto quando parliamo di dati e transazioni, e il monitoraggio delle modifiche ai dati non resta in disparte. Vediamo cosa potrebbe andare storto e cosa succederà con Debezium in ciascun caso. Ci sono tre opzioni di guasto: . Tuttavia, ci sono eccezioni. Se il connettore è stato disattivato per un lungo periodo (o non è riuscito a raggiungere l'istanza di MongoDB) e nel frattempo si è verificata la rotazione dell'oplog, al ripristino della connessione, il connettore riprenderà tranquillamente a leggere i dati dalla prima posizione disponibile, il che significa che parte dei dati in Kafka non verrà persa. Debezium è stata la mia prima esperienza con i sistemi CDC e globalmente è stata molto positiva. Il progetto ha impressionato per il supporto ai principali DBMS, la semplicità di configurazione, il supporto per il clustering e un'attiva comunità. Suggerisco a chi è interessato di dare un'occhiata alle 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 ottenere dati con una latenza minima. Il connettore JDBC (nella fornitura di Kafka Connect) esegue interrogazioni sulla tabella monitorata a intervalli fissi e (per lo stesso motivo) non genera messaggi quando i dati vengono eliminati (come si possono richiedere dati che non esistono?). Per affrontare esigenze simili, si possono considerare le seguenti soluzioni (oltre a Debezium): Leggete anche nel nostro blog: Fonte: habr.comtrasformazioni si verifica quanto segue: tutti gli eventi CDC dal database monitorato andranno a finire nell'argomento con il nome data.cdc.dbname. In caso contrario (senza queste impostazioni), Debezium creerebbe per impostazione predefinita un topic 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
}
}
INSERT: il valore precedente (before) è uguale null, e dopo è la stringa che è stata inserita. UPDATE: in payload.before mostra il precedente stato della stringa, mentre in payload.after c'è il nuovo stato con il merito 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"
}
}trasformazioni questa volta si fa quanto segue: si trasforma il nome del topic di destinazione dal schema .. in data.cdc.mongo_.Affidabilità
Conclusione
P.S.
