
Në punën time përballlem shpesh me zgjidhje teknike të reja/produkte programore, për të cilat informacioni në internetin shqip është mjaft i kufizuar. Me këtë artikull do të përpiqem të plotësoj një hapësirë të tillë me një shembull nga praktika ime e fundit, kur ishte e nevojshme të konfiguronte dërgimin e ngjarjeve CDC nga dy DB popullore (PostgreSQL dhe MongoDB) në klasterin Kafka me ndihmën e Debezium. Shpresoj që ky artikull përmbledhës, i cili u shfaq si rezultat i punës së kryer, do të jetë i dobishëm edhe për të tjerë.
ĂfarĂ« Ă«shtĂ« Debezium dhe nĂ« pĂ«rgjithĂ«si CDC?
â Ă«shtĂ« njĂ« pĂ«rfaqĂ«sues i kategorisĂ« sĂ« softuerit CDC (), dhe nĂ«se do tĂ« ishim mĂ« tĂ« saktĂ« â Ă«shtĂ« njĂ« grup konneksionesh pĂ«r DB tĂ« ndryshme, tĂ« pĂ«rshtatshme me kornizĂ«n Apache Kafka Connect.
Ky i cili përdor licencën Apache License v2.0 dhe sponsorohet nga kompania Red Hat. Zhvillimi është duke u kryer që nga viti 2016, dhe deri më tani ka mbështetje zyrtare për DB të mëposhtme: MySQL, PostgreSQL, MongoDB, SQL Server. Ekzistojnë gjithashtu konneksione për Cassandra dhe Oracle, por aktualisht ato janë në statusin "në akses të hershëm", dhe lëshimet e reja nuk garantojnë përputhshmërinë mbrapa.
Nëse e krahasojmë CDC me qasjen tradicionale (kur aplikacioni lexon të dhënat direkt nga DB), disa nga përparësitë kryesore janë realizimi i transmetimit të ndryshimeve të të dhënave në nivelin e rreshtave me vonesë të ulët, besueshmëri të lartë dhe disponibilitet. Këto dy pika të fundit arrijnë me anë të përdorimit të klasterit Kafka si magazinë ngjarjesh CDC.
Një tjetër avantazh është fakti që për ruajtjen e ngjarjeve përdoret një model i vetëm, prandaj aplikacioni përfundimtar nuk do të ketë nevojë të shqetësohet për nuancat e shfrytëzimit të DB-ve të ndryshme.
Së fundi, falë përdorimit të brokerit të mesazheve, hapet një mundësi për shkallëzimin horizontal të aplikacioneve që ndjekin ndryshimet në të dhëna. Në këtë rast, ndikimi në burimin e të dhënave reduktohet në minimum, pasi marrja e të dhënave nuk ndodh drejtpërdrejt nga DB, por nga klasteri Kafka.
Rreth arkitekturës së Debezium
Përdorimi i Debezium përmbledhet në një skemë kaq të thjeshtë:
DB (si burim tĂ« dhĂ«nash) â konnektor nĂ« Kafka Connect â Apache Kafka â konsumator
Si ilustrim, do të jap një skemë nga faqja e projektit:

Megjithatë, kjo skemë nuk më pëlqen shumë, pasi krijon përshtypjen se është e mundur vetëm përdorimi i konnektorit sink.
NĂ« tĂ« vĂ«rtetĂ«, situata Ă«shtĂ« ndryshe: mbushja e Data Lake tuaj (elementi i fundit nĂ« skemĂ«n e mĂ«sipĂ«rme) Ââ kjo nuk Ă«shtĂ« mĂ«nyra e vetme pĂ«r tĂ« pĂ«rdorur Debezium. Ngjarjet e dĂ«rguara nĂ« Apache Kafka mund tĂ« pĂ«rdoren nga aplikacionet tuaja pĂ«r tĂ« zgjidhur situata tĂ« ndryshme. PĂ«r shembull:
- heqja e të dhënave të papërshtatshme nga cache;
- dërgimi i njoftimeve;
- aktualizime të indekseve të kërkimit;
- një formë e diasporës së auditi;
- âŠ
Në rast se keni një aplikacion në Java dhe nuk keni nevojë/mundësi të përdorni një klaster Kafka, ka gjithashtu mundësinë e punës përmes . Avantazhi i qartë është se me të mund të hiqni dorë nga infrastruktura shtesë (si konektori dhe Kafka). Megjithatë, ky zgjidhje është shpallur e vjetruar (deprecated) që nga versioni 1.1 dhe nuk rekomandohet më për përdorim (në versionet e ardhshme mbështetje e tij mund të hiqet).
Në këtë artikull do të shqyrtohet arkitektura e rekomanduar nga zhvilluesit, e cila siguron qëndrushmëri dhe mundësi për shkallëzim.
Konfigurimi i konektorit
PĂ«r tĂ« filluar tĂ« ndjekim ndryshimet e vlerĂ«s mĂ« tĂ« rĂ«ndĂ«sishme â tĂ« dhĂ«nat, na nevojiten:
- burimi i të dhënave, që mund të jetë MySQL nga versioni 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- klasteri Apache Kafka;
- instanca Kafka Connect (versionet 1.x, 2.x);
- konektori i konfiguruar Debezium.
Puna për dy piketat e para, pra procesi i instalimit të DBMS dhe Apache Kafka, kalon kufijtë e këtij artikulli. Megjithatë, për ata që dëshirojnë të zhvillojnë gjithçka në një rreth të mbyllur, në depozitat zyrtare me shembuj ka një .
Ne do të ndalemi më shumë në dy pikat e fundit.
0. Kafka Connect
Këtu dhe më tej në artikull, të gjitha shembujt e konfigurimit shqyrtohen në kontekstin e imazhit Docker, që shpërndahen nga zhvilluesit e Debezium. Ai përmban të gjitha skedaret e nevojshme të plugjeve (konektorëve) dhe parashikon konfigurimin e Kafka Connect përmes variablave të mjedisit.
Në rast se planifikohet përdorimi i Kafka Connect nga Confluent, do të nevojitet të shtoni manualisht plugjet e konektorëve të nevojshëm në direktorinë e specifikuar në plugin.path ose e caktuar përmes variablave të mjedisit CLASSPATH. Cilësimet e punëtorit Kafka Connect dhe konektorëve përcaktohen përmes skedareve të konfigurimit, të cilat kalojnë si argumente në komandën e fillimit të punëtorit. Më shumë detaje shih në .
Procesi i konfigurimit të Debezium me konektorin realizohet në dy faza. Le të shqyrtojmë secilën prej tyre:
1. Konfigurimi i kuadrit Kafka Connect
Për transmetimin e të dhënave në klasterin Apache Kafka, në kuadrin Kafka Connect përcaktohen parametra specifikë, të tillë si:
- parametrat e lidhjes me klasterin,
- emrat e temave, në të cilat do të ruhet përveç kësaj konfigurimi i vetë konektorit,
- emri i grupit në të cilin është aktivizuar konektori (në rastin e përdorimit të modit të shpërndarë).
Imazhi zyrtar Docker i projektit mbĂ«shtet konfigurimin pĂ«rmes variablave tĂ« mjedisit â dhe kĂ«tĂ« do tĂ« pĂ«rdorim. Prandaj, shkarkojmĂ« imazhin:
docker pull debezium/connectGrupi minimal i variablave të mjedisit të nevojshëm për aktivizimin e konektorit është si më poshtë:
-
BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092â lista fillestare e serverĂ«ve tĂ« klasterit Kafka pĂ«r tĂ« marrĂ« listĂ«n e plotĂ« tĂ« anĂ«tarĂ«ve tĂ« klasterit; -
OFFSET_STORAGE_TOPIC=connector-offsetsâ tema pĂ«r ruajtjen e pozicioneve ku ndodhet aktualisht konektori; -
CONNECT_STATUS_STORAGE_TOPIC=connector-statusâ tema pĂ«r ruajtjen e statusit tĂ« konektorit dhe detyrave tĂ« tij; -
CONFIG_STORAGE_TOPIC=connector-configâ tema pĂ«r ruajtjen e tĂ« dhĂ«nave tĂ« konfigurimit tĂ« konektorit dhe detyrave tĂ« tij; -
GROUP_ID=1â identifikuesi i grupit tĂ« punĂ«torĂ«ve, nĂ« tĂ« cilĂ«t mund tĂ« ekzekutohet detyra e konektorit; i nevojshĂ«m gjatĂ« pĂ«rdorimit tĂ« modit tĂ« shpĂ«rndarĂ« (distributed) modi.
Aktivizoni kontejnerin me këto variabla:
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.2Shënim mbi Avro
Me default, Debezium shkruan të dhëna në formatin JSON, që është e pranueshme për ambiente provuese dhe volumet e vogla të të dhënave, por mund të bëhet problem në bazat e të dhënave me ngarkesë të lartë. Një alternativë për konvertuesin JSON është serializimi i mesazheve nëpërmjet formatsit bimor, që lejon uljen e ngarkesës në nën-sistemin I/O në Apache Kafka.
Për të përdorur Avro kërkohet të vendoset një (për ruajtjen e skemave). Variablat për konvertuesin do të duken si më poshtë:
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.AvroConverterDetajet pĂ«r pĂ«rdorimin e Avro dhe konfigurimin e regjistrit dalin jashtĂ« pĂ«rmbajtjes sĂ« kĂ«tij artikulli â mĂ« tej pĂ«r qartĂ«si ne do tĂ« pĂ«rdorim JSON.
2. Konfigurimi i konektorit
Tani mund të kalojmë direkt në konfigurimin e konektorit, i cili do të lexojë të dhënat nga burimi.
Le të shohim me shembujt e konektorëve për dy DBMS: PostgreSQL dhe MongoDB, - për të cilat kam përvojë dhe ku ka dallime (edhe pse të vogla, në disa raste - të rëndësishme!).
Konfigurimi përshkruhet në notacionin JSON dhe ngarkohet në Kafka Connect nëpërmjet një kërkese POST.
2.1. PostgreSQL
Shembulli i konfigurimit të konektorit për 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"
}
}Principi i funksionimit të konektorit pas një konfigurimi të tillë është mjaft i thjeshtë:
- Në nisjen e parë, ai lidhet me databazën e specifikuar në konfigurim dhe niset në modin snapshot fillestar, duke dërguar në Kafka një grup të dhënash fillestar të marra me
SELECT * FROM table_name. - Pasi të përfundojë inicializimi, konektori kalon në modin e leximit të ndryshimeve nga skedarët WAL të PostgreSQL.
Rreth opsioneve të përdorura:
-
emriâ emri i konektorit pĂ«r tĂ« cilin pĂ«rdoret konfigurimi i pĂ«rshkruar mĂ« poshtĂ«; mĂ« vonĂ« ky emĂ«r pĂ«rdoret pĂ«r tĂ« punuar me konektorin (dmth. tĂ« kontrollosh statusin/rishtazosh/rivendosĂ«sh konfigurimin) pĂ«rmes API REST tĂ« Kafka Connect; -
connector.classâ klasa e konektorit DBMS qĂ« do tĂ« pĂ«rdoret nga konektori qĂ« po konfigurohet; -
plugin.nameâ emri i plugin-it pĂ«r dekodimin logjik tĂ« tĂ« dhĂ«nave nga skedarĂ«t WAL. NĂ« dispozitĂ« janĂ«wal2json,decoderbuffsdhepgoutput. Dy tĂ« parat kĂ«rkojnĂ« instalimin e zgjatjeve pĂ«rkatĂ«se nĂ« DBMS, ndĂ«rsapgoutputpĂ«r PostgreSQL versionin 10 e lart nuk kĂ«rkon veprime tĂ« tjera; -
database.*â opsionet pĂ«r t'u lidhur me DB, kudatabase.server.nameâ emri i instancĂ«s PostgreSQL, i pĂ«rdorur pĂ«r tĂ« formuar emrin e temĂ«s nĂ« klasterin Kafka; -
table.include.listâ lista e tabelave nĂ« tĂ« cilat dĂ«shirojmĂ« tĂ« ndjekim ndryshimet; pĂ«rcaktohet nĂ« formatinschema.table_name; nuk mund tĂ« pĂ«rdoret sĂ« bashku metable.exclude.list; -
heartbeat.interval.msâ intervali (nĂ« milisekonda), me tĂ« cilin konektori dĂ«rgon mesazhe heartbeat nĂ« njĂ« temĂ« tĂ« veçantĂ«; -
heartbeat.action.queryâ kĂ«rkesa qĂ« do tĂ« ekzekutohet gjatĂ« dĂ«rgimit tĂ« çdo mesazhi heartbeat (opcija e re nga versioni 1.1); -
slot.nameâ emri i slotit tĂ« replikimit, qĂ« do tĂ« pĂ«rdoret nga konektori; publication.nameâ emri nĂ« PostgreSQL, qĂ« pĂ«rdor konektori. NĂ« rast se nuk ekziston, Debezium do tĂ« pĂ«rpiqet ta krijojĂ« atĂ«. NĂ«se pĂ«rdoruesi, me tĂ« cilin po lidhet, nuk ka mjaftueshĂ«m tĂ« drejta pĂ«r kĂ«tĂ« veprim â konektori do tĂ« ndalojĂ« me njĂ« error;-
transformonpërcakton se si të ndryshohet emri i temës së synuar:-
transforms.AddPrefix.typetregon që do të përdoren shprehje të rregullta; -
transforms.AddPrefix.regexâ maska sipas sĂ« cilĂ«s riemĂ«rohet emri i temĂ«s sĂ« synuar; -
transforms.AddPrefix.replacementâ konkretisht ajo nĂ« tĂ« cilĂ«n po riemĂ«rohet.
-
Më shumë për heartbeat dhe transforms
Në mënyrë të paracaktuar, konektori dërgon të dhëna në Kafka për çdo transaksion të komituar, dhe numri i tij LSN (Log Sequence Number) regjistrohet në një temë shërbimi offset. Por çfarë do të ndodhi nëse konektori është i konfiguruar për të lexuar jo të gjithë bazën e të dhënave, por vetëm një pjesë të tabelave të saj (ku përditësimi i të dhënave ndodh jo shpesh)?
- Konektori do të lexojë skedaret WAL dhe nuk do të zbulojë atje komitë për transaksionet në ato tabela, për të cilat po vëzhgon.
- Prandaj ai nuk do të përditësojë pozita e tij aktuale as në temë, as në slotin e replikimit.
- Kjo, nga ana tjetër, do të çojë në "mbajtjen" e skedarëve WAL në disk dhe ndoshta në shterjen e gjithë hapësirës diskale.
Dhe këtu ndihmojnë opsionet heartbeat.interval.ms dhe heartbeat.action.query. Përdorimi i këtyre opsioneve si çift lejon që çdo herë kur dërgohet një mesazh heartbeat të bëhet një kërkesë për të ndryshuar të dhënat në një tabelë të veçantë. Kjo mënyrë e mbajti gjithmonë të aktualizuar LSN, në të cilin ndodhet tani konektori (në slotin e replikimit). Kjo lejon që DBMS të fshijë skedarët WAL, të cilët nuk janë më të nevojshëm. Mësoni më shumë rreth funksionit të opsioneve në .
Një tjetër opsion, që meriton më shumë vëmendje, është transformon. Megjithatë, ajo është më shumë në lidhje me lehtësinë dhe estetikën...
Në mënyrë të paracaktuar, Debezium krijon tema duke u bazuar në politikën e mëposhtme të emërtimit: serverName.schemaName.tableName. Kjo nuk është gjithmonë e përshtatshme. Opsionet transformon mundë të përcaktoni listën e tabelave me anë të shprehjeve të rregullta, ngjarjet e të cilave duhet të rrugëzohen në një topic me emrin përkatës.
NĂ« konfigurimin tonĂ« falĂ« NĂ« pĂ«rfundim tĂ« pĂ«rshkrimit tĂ« konfigurimit tĂ« konektorit pĂ«r PostgreSQL, duhet tĂ« flasim pĂ«r karakteristikat/kufizimet e mĂ«poshtme tĂ« punĂ«s sĂ« tij: Pra, le tĂ« ngarkojmĂ« konfigurimin tonĂ« nĂ« konektor: KontrollojmĂ« qĂ« ngarkimi ka kaluar me sukses dhe konektori ka nisur: ShkĂ«lqyeshĂ«m: ai Ă«shtĂ« konfiguruar dhe gati pĂ«r punĂ«. Tani le tĂ« pretendonim tĂ« ishim konsumator dhe tĂ« lidhenim me Kafka, e pastaj tĂ« shtonim dhe ndryshonim njĂ« tĂ« dhĂ«nĂ« nĂ« tabelĂ«: NĂ« topicin tonĂ« kjo do tĂ« paraqitet si mĂ« poshtĂ«: NjĂ« JSON shumĂ« i gjatĂ« me ndryshimet tona NĂ« tĂ« dy rastet, regjistrimet pĂ«rbĂ«hen nga çelĂ«si (PK) i regjistrimit qĂ« Ă«shtĂ« ndryshuar dhe nĂ« vetĂ« thelbin e ndryshimeve: si ishte regjistrimi para dhe si u bĂ« pas. Ky konektor pĂ«rdor mekanizmin standard tĂ« replikimit tĂ« MongoDB, duke lexuar informacionin nga oplog-i i nodit kryesor tĂ« DB. NjĂ«lloj si konektori i pĂ«rshkruar mĂ« parĂ« pĂ«r PgSQL, kĂ«tu gjithashtu nĂ« fillim tĂ« instalimit merret njĂ« snapshot fillestar i tĂ« dhĂ«nave, pas sĂ« cilĂ«s konektori kalon nĂ« modin e leximit tĂ« oplog-ut. Shembulli i konfigurimit: Siç mund tĂ« vĂ«reni, kĂ«tu nuk ka opsionesh tĂ« reja krahasuar me shembullin e kaluar, por Ă«shtĂ« reduktuar vetĂ«m numri i opsioneve qĂ« lidhen me lidhjen me DB dhe prefikset e tyre. CilĂ«simet ĂĂ«shtja e qĂ«ndrueshmĂ«risĂ« dhe disponueshmĂ«risĂ« sĂ« lartĂ« nĂ« kohĂ«n tonĂ« Ă«shtĂ« mĂ« e rĂ«ndĂ«sishme se kurrĂ« â sidomos kur flasim pĂ«r tĂ« dhĂ«nat dhe transaksionet, dhe ndjekja e ndryshimeve tĂ« tĂ« dhĂ«nave nuk qĂ«ndron nĂ« anĂ«n e kĂ«tij problemi. Le tĂ« shqyrtojmĂ« se çfarĂ« mund tĂ« shkojĂ« keq dhe çfarĂ« do tĂ« ndodhte me Debezium nĂ« çdo rast. Ka tre variante tĂ« dĂ«shtimit: MegjithatĂ«, ka pĂ«rjashtime. NĂ«se konektori ka qenĂ« nĂ« gjendje tĂ« fikur pĂ«r njĂ« periudhĂ« tĂ« gjatĂ« (ose nuk ka mundur tĂ« arrijĂ« nĂ« instancĂ«n e MongoDB), dhe oplog-u ka kaluar rrotullimin gjatĂ« kĂ«tij kohĂ«, nĂ« rikthimin e lidhjes konektori do tĂ« vazhdojĂ« tĂ« lexojĂ« tĂ« dhĂ«nat nga pozita e parĂ« tĂ« disponueshme, duke rezultuar nĂ« faktin se njĂ« pjesĂ« e tĂ« dhĂ«nave nĂ« Kafka jo do tĂ« pĂ«rfshihet. Debezium â Ă«shtĂ« pĂ«rvoja ime e parĂ« me sistemet e CDC dhe nĂ« pĂ«rgjithĂ«si shumĂ« pozitive. Projekti Ă«shtĂ« tĂ«rheqĂ«s pĂ«r mbĂ«shtetje tĂ« DBMS-ve kryesore, thjeshtĂ«sinĂ« e konfigurimit, mbĂ«shtetje pĂ«r klasterizimin dhe njĂ« komunitet aktiv. PĂ«r ata qĂ« janĂ« tĂ« interesuar, rekomandoj tĂ« shikoni udhĂ«zuesit pĂ«r dhe . NĂ« krahasim me konektorin JDBC pĂ«r Kafka Connect, pĂ«rparĂ«sia kryesore e Debezium Ă«shtĂ« se ndryshimet lexohen nga regjistrat e DBMS, duke mundĂ«suar marrjen e tĂ« dhĂ«nave me vonesĂ« minimale. Konektori JDBC (nga paketa e Kafka Connect) bĂ«n kĂ«rkesa nĂ« tabelĂ«n e ndjekur me njĂ« interval tĂ« fiksuar dhe (pĂ«r kĂ«tĂ« arsye) nuk gjeneron mesazhe kur tĂ« dhĂ«nat fshihen (si mund tĂ« kĂ«rkosh tĂ« dhĂ«na qĂ« nuk ekzistojnĂ«?). PĂ«r tĂ« zgjidhur detyra tĂ« ngjashme, mund tĂ« shqyrtoni zgjidhjet e mĂ«poshtme (pĂ«rveç Debezium): Lexoni gjithashtu nĂ« blogun tonĂ«: Burimi: habr.comtransformon ndodh njĂ« e tillĂ«: tĂ« gjitha ngjarjet CDC nga baza e tĂ« dhĂ«nave tĂ« ndjekura do tĂ« shkojnĂ« nĂ« njĂ« topic me emrin data.cdc.dbname. PĂ«rndryshe (pa kĂ«to ndĂ«rrime) Debezium do tĂ« krijonte pĂ«r çdo tabelĂ« njĂ« topic tĂ« tipit: pg-dev.public..
Kufizimet e konektorit
Përdorimi i konfigurimit
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
}
}
SHTO: vlera para (para) Ă«shtĂ« e barabartĂ« me null, dhe pas â vargu qĂ« u fut. UPDATE: nĂ« payload.before tregon gjendjen e mĂ«parshme tĂ« rreshtit, ndĂ«rsa nĂ« payload.after â e re me thelbin e ndryshimeve.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"
}
}transformon në këtë rast bëjnë të siguiente: shndërrojnë emrin e temës që drejton në sistemin e .. në data.cdc.mongo_.Qëndrueshmëri
Përfundim
P.S.
