
Oma töös kohtan ma tihti uusi tehnilisi lahendusi/programmide tooteid, mille kohta on vene keeles internetis suhteliselt vĂ€he teavet. Selle artikliga pĂŒĂŒan tĂ€ita ĂŒhe sellise lĂŒnga, tuues nĂ€ite oma hiljutisest praktikast, kus oli vajalik seadistada CDC-sĂŒndmuste saatmine kahest populaarsest andmebaasist (PostgreSQL ja MongoDB) Kafka klastrisse Debeziumi abil. Loodan, et see ĂŒlevaateartikkel, mis sĂŒndis tehtud töö tulemusena, osutub kasulikuks ka teistele.
Mis on Debezium ja ĂŒldiselt CDC?
â CDC tarkvara kategooria esindaja (), tĂ€psemalt öeldes on see erinevate andmebaaside konnektorite kogum, mis on ĂŒhilduv Apache Kafka Connect raamistiku sĂŒsteemiga.
See millel on Apache License v2.0 litsents ja mida sponsoreerib ettevĂ”te Red Hat. Arendustööd on alustatud 2016. aastal ning praegu on ametlik tugi jĂ€rgmistele andmebaasidele: MySQL, PostgreSQL, MongoDB, SQL Server. Samuti on olemas konnektorid Cassandra ja Oracle jaoks, kuid need on praegu 'varajase juurdepÀÀsu' staatuses, ning uued vĂ€ljalasked ei garanteeri tagasipöördumatut ĂŒhilduvust.
Kui vĂ”rrelda CDC-d traditsioonilise lĂ€henemisega (kui rakendus loeb andmeid otse andmebaasist), siis selle peamisteks eelisteks on madala latentsusega, kĂ”rge usaldusvÀÀrsuse ja kĂ€ttesaadavuse taseme saavutamine andmete muutuste voogude realiseerimise kaudu ridade tasandil. Viimane kahest punktist saavutatakse, kasutades Kafka klastrit CDC-sĂŒndmuste salvestamiseks.
Samuti on eeliste hulka arvestatav, et sĂŒndmuste salvestamiseks kasutatakse ĂŒhtset mudelit, seega ei pea lĂ”pprakendus muretsema erinevate andmebaaside haldamise spetsiifikaga.
LÔpuks, thanks to the use of a message broker, the potential for horizontal scalability emerges for applications that track data changes. At the same time, the impact on the data source is minimized, as data is not retrieved directly from the database but rather from the Kafka cluster.
Debeziumi arhitektuurist
Debeziumi kasutamine piirdub sellise lihtsa skeemiga:
Andmebaas (andmeallikas) â konnektor Kafka Connectis â Apache Kafka â tarbija
Illustreerimiseks toome vÀlja skeemi projekti veebisaidilt:

Kuid see skeem ei meeldi mulle eriti, kuna jÀÀb mulje, et on vÔimalik kasutada ainult sink-konnektorit.
Tegelikult on olukord teine: teie Data Lake'i tĂ€itmine (ĂŒlemises skeemis viimane element) â pole ainus viis Debeziumi rakendamiseks. Apache Kafka'sse saadetud sĂŒndmusi saab teie rakendustes kasutada erinevate olukordade lahendamiseks. NĂ€iteks:
- mitteaktuaalsete andmete eemaldamine vahemÀlust;
- teadete saatmine;
- otsinguindeksite vÀrskendamine;
- mingisugused auditilogid;
- âŠ
Kui teie rakendus on Java-s ja puudub vajadus/vÔimalus kasutada Kafka klastrit, on olemas ka vÔimalus töötada lÀbi . Selge eelis on see, et sellega ei pea lisainfrastruktuurist (konnektor ja Kafka) loobuma. Siiski on see lahendus 1.1 versioonist alates kuulutatud aegunuks (deprecated) ja selle kasutamist ei soovitata enam (tulevastes vÀljaannetes vÔib selle toe eemaldada).
KÀesolevas artiklis kÀsitletakse soovitatud arendajate arhitektuuri, mis tagab talitlushÀireteta toimimise ja skaleeritavuse.
Ăhenduse konfiguratsioon
Et alustada muutuste jĂ€lgimist meie peamise varanduse â andmete â osas, on meil vaja:
- andmeallikat, milleks vÔivad olla MySQL alates versioonist 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- Apache Kafka klaster;
- Kafka Connecti instants (versioonid 1.x, 2.x);
- konfigureeritud Debeziumi ĂŒhendaja.
Kaks esimest punkti, st andmebaasi ja Apache Kafka installimisprotsess, jÀÀvad artikli huvist vÀljapoole. Siiski, neile, kes soovivad kÔik liivakasti kÀivitada, on ametlikus nÀidiste hoidlas olemas valmis .
Keskendume aga kahe viimase punkti sisule.
0. Kafka Connect
Siin ja edaspidi artiklis kaalutakse kĂ”iki konfigureerimise nĂ€iteid Debeziumi arendajate levitatava Docker-pildi kontekstis. See sisaldab kĂ”iki vajalikke pistikfailide (ĂŒhendajate) faile ja vĂ”imaldab Kafka Connecti konfigureerimist keskkonnamuutujate abil.
Kui eeldatakse, et kasutatakse Confluenti Kafka Connecti, tuleb vajalike ĂŒhendajate pistikfailid iseseisvalt lisada kausta, mis on mÀÀratud plugin.path vĂ”i mÀÀratletud keskkonnamuutujaga CLASSPATH. Kafka Connecti töötlus ja konnektorite seadistused mÀÀratakse konfigureerimisfailide kaudu, mis antakse töötluse kĂ€ivitamise kĂ€su argumentidena. Rohkem teavet leiate .
Kogu Debeizumi seadistamisprotsess konnektoriga toimub kahes etapis. Vaatame igaĂŒht neist:
1. Kafka Connect raamistiku seadistamine
Andmete voogedastamiseks Apache Kafka klastrisse Kafka Connect raamistiku kaudu mÀÀratakse spetsiifilised parameetrid, nagu:
- ĂŒhenduse parameetrid klastriga,
- teemade nimed, kus salvestatakse otse konnektori konfiguratsioon,
- grupi nimi, milles konnektor töötab (juhul kui kasutatakse jaotatud reĆŸiimi).
Projektile ametlik Docker-pilt toetab konfigureerimist keskkonnamuutujate abil â seda me ka kasutame. Nii et laadime alla pildi:
docker pull debezium/connectMinimum keskkonnamuutujate kogum, mis on vajalik konnektori kÀivitamiseks, on jÀrgmine:
-
BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092â algne Kafka klastrite serverite loend, et saada tĂ€ispakkumine klastrite liikmetest; -
OFFSET_STORAGE_TOPIC=connector-offsetsâ teema, mis salvestab hetkel konnektori asukohad; -
CONNECT_STATUS_STORAGE_TOPIC=connector-statusâ konnektori ja tema ĂŒlesannete oleku ladustamise teema; -
CONFIG_STORAGE_TOPIC=connector-configâ konnektori ja tema ĂŒlesannete konfiguratsioonide andmete ladustamise teema; -
GROUP_ID=1â töötajate grupi identifikaator, kus konnektori ĂŒlesanne vĂ”ib toimuda; vajalik jaotatud (distributed) reĆŸiimi kasutamisel.
KĂ€ivitame konteineri nende muutujaitega:
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.2MĂ€rkus Avro kohta
Vaikimisi kirjutab Debezium andmeid JSON-formaadis, mis on sobiv liivakastide ja vĂ€ikeste andmemahtude jaoks, kuid vĂ”ib osutuda probleemseks kĂ”rge koormusega andmebaasides. JSON-konverteri alternatiiviks on sĂ”numite serialiseerimine binaarseks formaadiks, mis vĂ”imaldab vĂ€hendada I/O alamsĂŒsteemi koormust Apache Kafka-s.
Avro kasutamiseks on vajalik eraldi (skeemide hoidmiseks). Konverteri muutujaid kasutatakse jÀrgmiselt:
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.AvroConverterAvro kasutamise ja registreerimise seadistamise ĂŒksikasjad ĂŒletavad artikli piire â edaspidi kasutame selguse huvides JSON-i.
2. Konnektori seadistamine
NĂŒĂŒd saame minna otse konnektori seadistamise juurde, mis loeb andmeid allikast.
Vaatame kahe andmebaasi, PostgreSQLi ja MongoDB, konnektoreid â nendes on mul kogemusi ja on teatud erinevusi (kuigi vĂ€ikesed, vĂ”ivad need mĂ”nel juhul olla mĂ€rkimisvÀÀrsed!).
Konfiguratsioon on kirjeldatud JSON-i notatsioonis ja laaditakse Kafka Connecti POST-requests abil.
2.1. PostgreSQL
PostgreSQL konnektori konfiguratsiooni nÀide:
{
"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"
}
}Konnektori tööpĂ”himĂ”te pĂ€rast sellist seadistust on ĂŒsna lihtne:
- Esimese kĂ€ivitamise ajal ĂŒhendub see konfigureeritud andmebaasiga ja siseneb reĆŸiimi esialgne snapshot, saates Kafka'sse esmased andmeĂŒksused, mis on saadud tingimuslikust
SELECT * FROM table_name. - PĂ€rast initsialiseerimise lĂ”petamist liigub konnektor PostgreSQL WAL-failidest muudatuste lugemise reĆŸiimi.
Kasutatavate valikute kohta:
-
nameâ konnektori nimi, mille jaoks kasutatakse allpool kirjeldatud seadistust; tulevikus kasutatakse seda nime konnektoriga töötamiseks (nt oleku vaatamine/taaskĂ€ivitamine/seadistuse vĂ€rskendamine) Kafka Connect REST API kaudu; -
connector.classâ andmepood, mida konfigureeritav ĂŒhendus kasutab; -
plugin.nameâ pistiku nimi, mis on mĂ”eldud andmete loogiliseks dekodeerimiseks WAL-failidest. Saadaval on jĂ€rgmised valikud:wal2json,decoderbuffsjapgoutput. Esimene kaks vajavad vastava laienduse installimist andmebaasis, samas kuipgoutputPostgreSQL versiooni 10 ja uuemate jaoks ei vajata tĂ€iendavaid toiminguid; -
database.*â ĂŒhenduse valikud andmebaasiga, kusdatabase.server.nameâ PostgreSQL instantsi nimi, mida kasutatakse teema nime genereerimiseks Kafka klastris; -
table.include.listâ tabelite loetelu, kus soovime jĂ€lgida muudatusi; sisestatakse vormingusschema.table_name; ei saa kasutada koostable.exclude.list; -
heartbeat.interval.msâ intervall (millisekundites), mil konnektor saadab heartbeat-sĂ”numeid spetsiaalsesse teema; -
heartbeat.action.queryâ pĂ€ring, mida tehakse iga heartbeat-sĂ”numi saatmisel (valik, mis on saadaval alates versioonist 1.1); -
slot.nameâ replikatsiooni sloti nimi, mida konnektor kasutab; publication.nameâ nimi PostgreSQL-is, mida ĂŒhendus kasutab. Kui see puudub, proovib Debezium selle luua. Kui kasutajal, kelle all ĂŒhendus toimub, ei ole piisavalt Ă”igusi selle toimingu teostamiseks, lĂ”petab ĂŒhendus töö tĂ”rketeatega;-
muundamisedmÀÀrab, kuidas tÀpselt sihtteema nime muuta:-
transforms.AddPrefix.typenÀitab, et kasutame regulaaravaldisi; -
transforms.AddPrefix.regexâ muster, mille jĂ€rgi muudame sihtteema nime; -
transforms.AddPrefix.replacementâ see, millele me selle muutume.
-
Rohkem teavet heartbeat'i ja muundamiste kohta
Vaikimisi saadab ĂŒhendus andmed Kafka'sse iga kinnitatud tehingu korral ning selle LSN (Log Sequence Number) salvestatakse teenuse teema offset. Aga mis juhtub, kui ĂŒhendus on seadistatud lugema mitte kogu andmebaasi, vaid ainult osa selle tabelitest (kus andmete vĂ€rskendamine ei toimu sageli)?
- Ăhendus hakkab lugema WAL-faile ega leia neist tehingute kinnitamist nendes tabelites, mille ĂŒle ta jĂ€lgib.
- SeetÔttu ei uuenda ta oma praegust positsiooni ei teemas ega replikatsiooni slotis.
- See toimetab, et WAL-failid jÀÀvad kettale âkinniâ ja kogu kettaruumi ammendamine vĂ”ib juhtuda.
Siinkohal tulevad appi valikud heartbeat.interval.ms ja heartbeat.action.query. Nende valikute kasutamine koos vĂ”imaldab iga kord, kui saadetakse sĂŒdamepekslemise sĂ”num, teha pĂ€ring andmete muutmiseks eraldi tabelis. Sellisel viisil uuendatakse pidevalt LSN-i, kus konnektor praegu asub (replikatsiooni pesas). See vĂ”imaldab DBMS-il eemaldada WAL-failid, mis ei ole enam vajalikud. Lisainfot valikute töö kohta leiate .
Teine valik, mis vÀÀrib suuremat tÀhelepanu, on muundamised. Kuigi see on pigem mugavuse ja ilu pÀrast...
Vaikimisi loob Debezium teemasid, jĂ€rgides jĂ€rgmist nimede mÀÀramise poliitikat: serverName.schemaName.tableName. See ei pruugi alati mugav olla. Valikute abil muundamised vĂ”ib regulaarselt vĂ€ljendit kasutades mÀÀrata tabelite nimekirja, millelt sĂŒndmused tuleb suunata konkreetse nimega teema.
Meie konfiguratsioonis tĂ€nu PostgreSQL konnektori konfiguratsiooni kirjelduse lĂ”petamiseks tasub rÀÀkida jĂ€rgmistest omadustest/piirangutest selle toimimises: Kontrollime, et laadimine toimus edukalt ja konnektor kĂ€ivitus: Kontrollime, et laadimine Ă”nnestus ja konnektor kĂ€ivitati: SuurepĂ€rane: see on seadistatud ja valmis töötama. NĂŒĂŒd kĂ€itume tarbijana ja ĂŒhendame end Kafka'ga, seejĂ€rel lisame ja muudame tabelisse kirje: Meie teemas kuvab see jĂ€rgmiselt: VĂ€ga pikk JSON meie muudatustega MĂ”lemal juhul koosnevad kirjed jĂ€rjekorra (PK) vĂ”tmetest, mis on muudetud, ja muudatuste sisust: millisena oli kande sisu enne ja millisena ta pĂ€rast muutus. See konnektor kasutab MongoDB standardset replikatsiooni mehhanismi, lugedes teavet pĂ”hivĂ”lvi oplog'ist. Sarnaselt juba kirjeldatud PgSQL konnektorile, tehakse siin esmakordsel kĂ€ivitamisel andmete esmane snapshot, pĂ€rast mida vahetab konnektor ĂŒle oplog'i lugemise reĆŸiimile. Konfiguratsiooni nĂ€idis: Nagu mĂ€rkida, ei ole siin uusi vĂ”imalusi vĂ”rreldes eelmise nĂ€itega, kuid on vĂ€henenud vaid andmebaasiĂŒhenduste ja nende prefikside arv. Seaded TĂ”rke- ja kĂ”rge kĂ€ttesaadavuse kĂŒsimus on tĂ€napĂ€eval aktuaalne, eriti kui rÀÀgime andmetest ja tehingutest, ning andmete muudatuste jĂ€lgimine ei jÀÀ selles osas ĂŒle. Vaatame, mis vĂ”ib pĂ”himĂ”tteliselt valesti minna ning mis juhtub Debeziumiga igas neist olukordadest. On kolm tĂ”rkevarianti: Kuid on erandeid. Kui konnektor on pikka aega olnud vĂ€ljas (vĂ”i ei saanud ĂŒhendust MongoDB eksemplariga), ja oplog on selle aja jooksul pöörlemise lĂ€binud, siis ĂŒhenduse taastamisel jĂ€tkab konnektor rahulikult andmete lugemist esimeselt kĂ€tte saadavalt positsioonilt, mistĂ”ttu osa andmeid jĂ”uab Kafka ei kuni. Debezium on minu esimene kogemus CDC-sĂŒsteemidega, ning ĂŒldiselt on see vĂ€ga positiivne. Projekt teeb mulje oma toetuse poolest peamistele DBMS-idele, konfigureerimise lihtsuse, klastritoe ja aktiivse kogukonna poolest. Praktikat huvitavatele soovitan tutvuda juhenditega ja . Debezium'i peamine eelis JDBC-ĂŒhendaja ees Kafka Connect'i jaoks on see, et muudatused loetakse andmebaasi logidest, mis vĂ”imaldab andmeid saada minimaalse viivitusega. JDBC Connector (Kafka Connect'i tarnekomplektist) teeb pĂ€ringuid jĂ€lgitava tabeli kohta kindla ajavahega ja (sama pĂ”hjusel) ei genereeri see sĂ”numeid andmete kustutamisel (kuidas saab kĂŒsida andmeid, mida enam pole?). Sarnaste probleemide lahendamiseks tasub vaadata jĂ€rgmisi lahendusi (Lisaks Debeziumile): Lugege ka meie blogist: Allikas: habr.commuundamised juhtuda jĂ€rgmine: kĂ”ik CDC-sĂŒndmused jĂ€lgitavast andmebaasist jĂ”uavad teema nimega data.cdc.dbnameKui need seaded pole paigas, siis Debezium loodi vaikimisi igale tabelile teema pĂ”hjal: pg-dev.public..
Konnektori piirangud
Seega laadime meie konfigureerimise konnektorisse:
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: vÀÀrtus enne (before) on vÔrreldav null, ja pÀrast on rida, mis on lisatud. KUUDA: payload.before kuvatakse rea eelmine seisund, samas kui payload.after on uus koos muudatuste sisuga.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"
}
}muundamised Seekord tehakse jĂ€rgmist: sihtteema nimi muudetakse skeemiks .. ĂŒhes data.cdc.mongo_.Veadeta töö
KokkuvÔte
P.S.
