
Oma töös kohtan sageli uusi tehnilisi lahendusi/softtooteid, mille kohta on venekeelses internetis ĂŒsna vĂ€he teavet. Selle artikliga pĂŒĂŒan tĂ€ita ĂŒht sellist puudujÀÀki, tuues nĂ€ite oma hiljutisest praktikast, kus tuli seadistada CDC-sĂŒndmuste saatmine kahest populaarset andmebaasist (PostgreSQL ja MongoDB) Kafka klastrisse Debeziumi abil. Loodan, et see ĂŒlevaade, mis on kirjutatud tehtud töö tulemuste pĂ”hjal, osutub kasulikuks ka teistele.
Mis on Debezium ja ĂŒldiselt CDC?
on CDC () tarkvara kategooria esindaja, tĂ€psemalt on see konnektorite kogum erinevatele andmebaasidele, mis on ĂŒhilduvad Apache Kafka Connect raamistikuga.
See on millel on Apache License v2.0 litsents ja mille sponsoreerib Red Hat. Arendustööd on alustatud 2016. aastal ja praeguseks on selles ametlik tugi jĂ€rgmistele andmebaasidele: MySQL, PostgreSQL, MongoDB, SQL Server. Samuti on olemas konnektorid Cassandra ja Oracle jaoks, kuid praegu on need "varajase juurdepÀÀsu" staatuses ning uued vĂ€ljaanded ei taga tagasipöörduvat ĂŒhilduvust.
CDC-d vĂ”rreldes traditsioonilise lĂ€henemisega (kui rakendus loeb andmeid otse andmebaasist) on selle peamisteks eelisteks andmete muutuste streamimise realiseerimine ridade tasemel koos madala latentsuse, kĂ”rge usaldusvÀÀrsuse ja kĂ€ttesaadavusega. Viimane kaks punkti saavutatakse, kasutades CDC-sĂŒndmuste hoidmiseks Kafka klastrit.
Sarnaselt on eelisteks ka see, et sĂŒndmuste salvestamiseks kasutatakse ĂŒhtset mudelit, mistĂ”ttu lĂ”pp-rakendusel ei pea olema mures erinevate andmebaaside kasutamise nĂŒansside ĂŒle.
L finalmente, tÀnu sÔnumite vahendaja kasutamisele avardub rakenduste horisontaalne skaleerimine, mis jÀlgivad andmete muutusi. Samuti on andmeallika mÔju minimaalselt piiratud, kuna andmeid saadakse mitte otse andmebaasist, vaid Kafka klastrist.
Debeziumi arhitektuurist
Debeziumi kasutamine seisneb sellises lihtsas skeemis:
Andmebaas (kui andmeallikas) â konnektor Kafka Connectis â Apache Kafka â tarbija
Illustreerimiseks toome vÀlja skeemi projekti veebilehelt:

Kuid see skeem mulle ei meeldi, kuna mul on tunne, et on vÔimalik ainult sink-konektori kasutamine.
Kuid tegelikult on olukord teistsugune: teie Data Lake'i tĂ€itmine (viimane element ĂŒlaltoodud skeemis) â see ei ole ainus viis Debeziumi rakendamiseks. Apache Kafkasse edastatud sĂŒndmusi saavad teie rakendused kasutada erinevate olukordade lahendamiseks. NĂ€iteks:
- vÀheharjutatud andmete kustutamine vahemÀlust;
- teavituste saatmine;
- otsingumootorite indeksite vÀrskendamine;
- mingi sarnane auditi logi;
- âŠ
Kui teil on Java-rakendus ega ole vaja/vÔimalust kasutada Kafka klastrit, on samuti vÔimalik kasutada . Ilmselge pluss on see, et selle abil saab vÀltida lisainfrastruktuuri (nagu konektor ja Kafka) kasutamist. Kuid see lahendus on versioonist 1.1 deprekeeritud ja rohkem ei soovitata. Tulevastes vÀljaannetes vÔib selle toe eemaldada.
Selles artiklis kÀsitletakse soovitatavat arhitektuuri, mille on vÀlja töötanud arendajad, et tagada tÔrkeoleku kindlus ja skaleeritavus.
Konektori konfiguratsioon
Andmete, millel on kÔige suurem vÀÀrtus, muudatuste jÀlgimiseks, on meil vajalik:
- andmeallikas, milleks vÔivad olla MySQL versioonist 5.7 alates, PostgreSQL 9.6+, MongoDB 3.2+ ();
- Apache Kafka klaster;
- Kafka Connect instants (versioonid 1.x, 2.x);
- konfigureeritud Debeziumi konektor.
Esimese kahe punkti, st andmebaasi ja Apache Kafkade installimise protsess, ĂŒletab artikli piirid. Kuid neile, kes soovivad kĂ”ike liivakastis kĂ€ivitada, on ametlikus nĂ€idiste hoidlas valmis .
Keskendume aga kahele viimasele punktile.
0. Kafka Connect
Siin ja edaspidi artikkel kÔik konfiguratsiooni nÀited kÀsitletakse Debeziumi arendajate levitatud Docker-pildis. See sisaldab kÔiki vajalikke pluginifailisid (konektorid) ning eeldab Kafka Connect'i seadistamist keskkonnamuutujate kaudu.
Kui kavatsete kasutada Confluenti Kafka Connect'i, peate ise vajalikud konektorite pluginad lisama katalooge, nagu on nÀidatud plugin.path vÔi seadistada keskkonnamuutujaga CLASSPATH. Kafka Connecti ja konnektorite seaded mÀÀratakse konfiguratsioonifailide kaudu, mis edastatakse argumendina töötaja kÀivitamise kÀsule. Lisateavet vt .
Kogu Debeziumi seadistamise protsess konnektori variandiga toimub kahes etapis. Vaatame igaĂŒht neist:
1. Kafka Connecti raamistikku seadistamine
Andmete voogedastamiseks Apache Kafka klastrisse seadistatakse Kafka Connecti raamistikus spetsiifilised parameetrid, nagu:
- ĂŒhenduse parameetrid klastri juurde,
- teemade nimed, kus konnektori seadistused endid hoitakse,
- grupi nimi, milles konnektor on kĂ€ivitatud (jaotatud reĆŸiimi kasutamisel).
Projekti ametlik Docker-pilt toetab seadistamist keskkonna muutujate abil â seda me ka kasutame. Nii et, tĂ”mmake pilt:
docker pull debezium/connectMinimaalne muutujate komplekt, mis on vajalik konnektori kÀivitamiseks, nÀeb vÀlja jÀrgmine:
-
BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092â klastrisse Kafka ĂŒhendusserversite esialgne loend, et saada tĂ€ielik loend klastriliikmetest; -
OFFSET_STORAGE_TOPIC=connector-offsetsâ teema, kus hoitakse positsioone, kus konnektor hetkel asub; -
CONNECT_STATUS_STORAGE_TOPIC=connector-statusâ teema, kus hoitakse konnektori ja selle ĂŒlesannete staatust; -
CONFIG_STORAGE_TOPIC=connector-configâ teema, kus hoitakse konnektori ja selle ĂŒlesannete seadistuste andmeid; -
GROUP_ID=1â töötajate grupi ID, millel konnektori ĂŒlesanne vĂ”ib töötada; vajalik jaotatud (distributed) reĆŸiimi kasutamisel.
KĂ€ivita konteiner nende muutujatega:
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 andmed JSON-formaadis, mis sobib liivakastide ja vĂ€ikeste andmemahtude jaoks, kuid vĂ”ib muutuda probleemiks kĂ”rgelt koormatud andmebaasides. Alternatiiv JSON-konverterile on sĂ”numite serialiseerimine binaarformaati, mis vĂ”imaldab vĂ€hendada I/O alamsĂŒsteemi koormust Apache Kafka-s.
Avro kasutamiseks on vajalik eraldi (skeemide hoidmiseks). Muutujad konverteri jaoks nÀevad vÀlja jÀrgmised:
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 registri seadistamise ĂŒksikasjad jÀÀvad artikli raames vĂ€ljapoole â edaspidi kasutame selguse huvides JSON-i.
2. Konnektori seadistamine
NĂŒĂŒd saame liikuda otse konnektori konfiguratsiooni juurde, mis loeb andmeid allikast.
Vaatame kahe andmebaasi konnektori nĂ€itel: PostgreSQL ja MongoDB â mul on nende kohta kogemus ning nendes on erinevusi (kuigi vĂ€ikseid, vĂ”ivad need mĂ”nel juhul olla olulised!).
Konfiguratsioon on kirja pandud JSON-notatsioonis ja laaditakse Kafka Connecti POST-pÀringu kaudu.
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 seadistamist on ĂŒsna lihtne:
- Esimese kĂ€ivitamise korral ĂŒhendub see konfiguratsioonis mĂ€rgitud andmebaasi ja kĂ€ivitub reĆŸiimis initial snapshot, saates Kafka-sse algse andmekogumi, mis on saadud tingimusliku
SELECT * FROM table_name. - PĂ€rast initsialiseerimise lĂ”ppemist lĂ€heb konnektor ĂŒle muutuste lugemise reĆŸiimile PostgreSQL WAL-failidest.
Kasutatavatest valikutest:
-
nimiâ konnektori nimi, mille jaoks kasutatakse allpool kirjeldatud konfiguratsiooni; hiljem kasutatakse seda nime konnektoriga töötamiseks (st staatuse vaatamiseks/kĂ€ivitamiseks/konfiguratsiooni vĂ€rskendamiseks) lĂ€bi Kafka Connect REST API; -
connector.classâ andmebaasi konnektori klass, mida konfigureeritav konnektor kasutab; -
plugin.nameâ pluginanimi, mis on vajalik andmete loogiliseks dekodeerimiseks WAL-failidest. Saadaval on valikudwal2json,decoderbuffsjapgoutput. Esimesed kaks nĂ”uavad vastavate laienduste installimist andmebaasi, kuidpgoutputPostgreSQL versioon 10 ja uuem ei vaja tĂ€iendavaid toiminguid; -
database.*â valikud andmebaasi ĂŒhendamiseks, kusdatabase.server.nameâ PostgreSQL instantsi nimi, mida kasutatakse teema nime loomiseks Kafka klastris; -
table.include.listâ nimekiri tabelitest, kus soovime muudatusi jĂ€lgida; mÀÀratakse formaadisschema.table_name; ei tohi kasutada koostable.exclude.list; -
heartbeat.interval.msâ intervall (millisekundites), mil konnektor saadab heartbeat-sĂ”numeid spetsiaalsesse teema; -
heartbeat.action.queryâ pĂ€ring, mis tĂ€idetakse iga heartbeat-sĂ”numi saatmise ajal (valik ilmus versioonis 1.1); -
slot.nameâ replikatsioonislot'i nimi, mida konnektor kasutab; publication.nameâ nimi PostgreSQL-is, mida konnektor kasutab. Kui seda ei eksisteeri, proovib Debezium selle luua. Kui kasutajal, kelle alt ĂŒhendus toimub, ei ole selle toimingu jaoks piisavalt Ă”igusi, lĂ”petab konnektor töö veaga;-
transformsmÀÀrab, kuidas tÀpselt muuta sihteema nime:-
transforms.AddPrefix.typenÀitab, et kasutame regulaaravaldusi; -
transforms.AddPrefix.regexâ mask, mille alusel sihteema nimi mÀÀratakse uuesti; -
transforms.AddPrefix.replacementâ see, millele me selle ĂŒmber mÀÀrame.
-
Rohkem infot heartbeat'i ja transforms'i kohta
Vaikimisi saadab konnektor andmeid Kafka-sse igas komitendis tehingus ning selle LSN (Log Sequence Number) salvestatakse teenindusteemasse offset. Kuid mis juhtub, kui konnektor on seadistatud lugema mitte kogu andmebaasi, vaid ainult osa selle tabelitest (kus andmete vÀrskendamine ei toimu sageli)?
- Konnektor loeb WAL-faile ja ei leia neist tehingute komiteid nendesse tabelitesse, millele ta tÀhelepanu pöörab.
- SeetÔttu ei uuenda ta oma praegust positsiooni ei teemadesse ega replikatsioonislotis.
- See omakorda toob kaasa WAL-failide "hoidmise" kettal ja vÔimaliku kogu kettaruumi ammendumise.
Ja siin tulevad appi valikud heartbeat.interval.ms ja heartbeat.action.query. Nende valikute kooskasutamine vÔimaldab igal ajal heartbeat-sÔnumi saatmisel teostada pÀringu andmete muutmiseks eraldi tabelis. Sellega ajakohastatakse pidevalt LSN-i, kus konnektor praegu asub (replikatsioonislotis). See vÔimaldab andmebaasi halduril kustutada WAL-failid, mis enam ei ole vajalikud. Lisainfot nende valikute töötamisest saab .
Teine vĂ”imalus, mis vÀÀrib suuremat tĂ€helepanu, on transforms. Kuigi see on pigem mugavuse ja ilu kohtaâŠ
Debezium loob vaikimisi teemasid, jĂ€rgides jĂ€rgmisi nimetamisreegleid: serverName.schemaName.tableName. See ei pruugi alati mugav olla. VĂ”imalustega transforms saab kasutada regulaaravaldisi, et mÀÀrata tabelite loend, mille sĂŒndmusi peab suunama kindla nimega teema.
Meie konfiguratsioonis toimub tĂ€nu PostgreSQL konnektori konfiguratsiooni kokkuvĂ”tteks tasub rÀÀkida jĂ€rgmistest omadustest/piirangutest: Seega laadime meie konfiguratsiooni konnektorisse: Kontrollime, et laadimine lĂ€ks hĂ€sti ja konnektor kĂ€ivitati: SuurepĂ€rane: see on seadistatud ja tööks valmis. NĂŒĂŒd teeme nĂ€gu, et oleme tarbija ja ĂŒhendume Kafka'ga, seejĂ€rel lisame ja muudame kirjet tabelis: Meie teemas peegelduvad need jĂ€rgmise kujul: VĂ€ga pikk JSON meie muudatustega MĂ”lemal juhul koosnevad salvestused kirje (PK) vĂ”tmetest, mis on muudetud, ja muudatuste olemusest: milline oli kirje enne ja milline pĂ€rast. See konnektor kasutab MongoDB standardset replikatsiooni mehhanismi, lugedes teavet primaarse andmebaasi oplogâist. Sarnaselt eelnevalt kirjeldatud PgSQL konnektorile, teeb ka see esmakordsel kĂ€ivitamisel esialgse andmete snapshotsi, pĂ€rast mida lĂŒlitub konnektor oplog'i lugemise reĆŸiimile. Konfiguratsiooni nĂ€idis: Nagu mĂ€rgata, ei ole siin uusi vĂ”imalusi vĂ”rreldes eelmise nĂ€itega, kuid on vĂ€henenud vaid andmebaasi ĂŒhendamise ja nende eeltehtud nimetuste arv. Seaded Vastus vastupidavuse ja kĂ”rge k disponibiliteedi kĂŒsimusele on tĂ€napĂ€eval ÀÀrmiselt oluline â eriti kui rÀÀgime andmetest ja tehingutest, ning andmete muutuste jĂ€lgimine ei jÀÀ selles kĂŒsimuses kĂ”rvale. Vaatleme, mis vĂ”ib tegelikult valesti minna ja mis toimub Debeziumiga igas juhtumis. On kolm tĂ”rkevarianti: Kuid on ka erandeid. Kui ĂŒhendus on pikka aega olnud vĂ€lja lĂŒlitatud (vĂ”i ei ole saanud MongoDB instantsile ligi), ja oplog on selle aja jooksul pööratud, siis ĂŒhenduse taastumisel jĂ€tkab ĂŒhendus vaikselt andmete lugemist esimesest kĂ€ttesaadavast kohast, mille tĂ”ttu osa andmeid jĂ”uab Kafka. ei kaotsi. Debezium on minu esimene kogemus CDC-sĂŒsteemidega ja ĂŒldiselt vĂ€ga positiivne. Projekt paelus mind peamiste andmebaaside toega, lihtsa seadistamise, klastritoe ja aktiivse kogukonnaga. Huvi korral soovitan tutvuda juhenditega ja . Debeziumi peamine eelis Kafka Connect JDBC-ĂŒhendaja ees on see, et muudatused loetakse vĂ€lja andmebaasi pĂ€evikutest, mis vĂ”imaldab andmeid saada minimaalse viivitusega. JDBC-ĂŒhendaja (mis on osa Kafka Connectist) pĂ€rib jĂ€lgitavast tabelist kindla aja jooksul ja (sama pĂ”hjuse tĂ”ttu) ei genereeri teateid andmete kustutamisel (kuidas saab kĂŒsida andmeid, mida pole?). Sarnaste probleemide lahendamiseks tasub vaadata jĂ€rgmisi lahendusi (vĂ€lja arvatud Debezium): Lugege ka meie blogist: Allikas: habr.comtransforms jĂ€rgnev: kĂ”ik CDC-sĂŒndmused jĂ€lgitavast andmebaasist jĂ”uavad teema nimega data.cdc.dbname. Vastupidisel juhul (ilma nende seadistusteta) oleks Debezium vaikimisi loonud iga tabeli jaoks teema, mille formaat on: pg-dev.public..
Ăhenduse piirangud
Konfiguratsiooni rakendamine
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Ôrdne null, ja pÀrast on sisestatud string. UPDATE: payload.before kuvatakse rea eelmine olek, ja payload.after on uus koos muudatuste olemusega.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 seekord teevad nad jÀrgmist: muudavad sihttopiku nime mustrist .. ja data.cdc.mongo_.Talitluskatkestustunne
KokkuvÔte
P.S.
