
Në punën time, shpesh kam përballë zgjidhje teknike ose produkte softuerike të reja, për të cilat informacioni në internetin në gjuhën ruse është mjaft i kufizuar. Në këtë artikull, do të përpiqem të mbush një të tillë boshllëk me një shembull nga praktika ime e fundit, kur do duhej të konfiguronte dërgimin e ngjarjeve CDC nga dy DBMS të njohura (PostgreSQL dhe MongoDB) në klasterin Kafka përmes Debezium. Shpresoj që ky artikull përmbledhës, i cili është rezultat i punës së kryer, do të jetë i dobishëm edhe për të tjerët.
ĂfarĂ« Ă«shtĂ« Debezium dhe çfarĂ« Ă«shtĂ« CDC?
â pĂ«rfaqĂ«sues i kategorisĂ« sĂ« softuerit CDC (), dhe mĂ« saktĂ« â Ă«shtĂ« njĂ« koleksion konnektoresh pĂ«r DBMS tĂ« ndryshme, tĂ« pĂ«rshtatshme me kornizĂ«n Apache Kafka Connect.
Kjo i licencuar me Apache License v2.0 dhe i sponsorizuar nga kompania Red Hat. Zhvillimi është duke u bërë që nga viti 2016 dhe deri më tani mbështetur zyrtarisht janë këto DBMS: MySQL, PostgreSQL, MongoDB, SQL Server. Po ashtu ekzistojnë konnektoresh për Cassandra dhe Oracle, por aktualisht ata janë në statusin 'qasje të hershme', dhe lëshime të reja nuk garantojnë përshtatshmërinë e prapme.
Në krahasim me qasjen tradicionale (kur aplikacioni lexon të dhënat drejtpërdrejt nga DBMS), disa nga avantazhet kryesore të CDC përfshijnë implementimin e transmetimit të ndryshimeve të të dhënave në nivelin e rreshtave me vonesë të ulët, besueshmëri të lartë dhe disponueshmëri. Këto dy pika arrijnë përmes përdorimit të klasterit Kafka si ruajtje për ngjarjet e CDC.
Një tjetër avantaj është fakti që për ruajtjen e ngjarjeve përdoret një model i vetëm, kështu që aplikacioni përfundimtar nuk do të ketë nevojë të shqetësohet për nuancat e funksionimit të ndryshëm DBMS-ve.
Më në fund, falë përdorimit të brokerit të mesazheve, hapet mundësia për shkallëzim horizontal të aplikacioneve që ndjekin ndryshimet në të dhëna. Ndërkohë, ndikimi mbi burimin e të dhënave minimalizohet, pasi marrja e të dhënave ndodh jo drejtpërdrejt nga DBMS, por nga klasteri Kafka.
Për arkitekturën Debezium
Përdorimi i Debezium përmbledh një skemë kaq të thjeshtë:
DBMS (si burim tĂ« dhĂ«nash) â konektor nĂ« Kafka Connect â Apache Kafka â konsumator
Si ilustërim, do të jap një diagram nga faqja e projektit:

Megjithatë, kjo skemë nuk më pëlqen shumë, sepse krijon përshtypjen se është e kufizuar vetëm në përdorimin e sink-connector.
NĂ« tĂ« vĂ«rtetĂ«, situata Ă«shtĂ« ndryshe: mbushja e Data Lake tuaj (lidhja e fundit nĂ« skemĂ«n mĂ« sipĂ«r) â nuk Ă«shtĂ« mĂ«nyra e vetme e pĂ«rdorimit tĂ« 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ë panevojshme nga cache;
- dërgimi i njoftimeve;
- përditësimi i indekseve të kërkimit;
- njëfarë logu auditimi;
- âŠ
Në rast se keni një aplikacion në Java dhe nuk keni nevojë/mundësi për të përdorur një klaster Kafka, ka gjithashtu mundësinë e punës përmes . Avantazhi i qartë është se me të mund të heqësh dorë nga infrastruktura shtesë (si connector-i dhe Kafka). Megjithatë, kjo zgjidhje është shpallur e vjetruar (deprecated) që nga versioni 1.1 dhe nuk rekomandohet më për t'u përdorur (në lëshimet e ardhshme, mbështetja e saj mund të hiqet).
Në këtë artikull do të shqyrtohet arkitektura e rekomanduar nga zhvilluesit, e cila ofron qëndrueshmëri dhe mundësi për shkallëzim.
Konfigurimi i lidhësit
PĂ«r tĂ« filluar tĂ« ndjekim ndryshimet nĂ« vlerĂ«n mĂ« tĂ« rĂ«ndĂ«sishme â tĂ« dhĂ«nat, â na nevojiten:
- burimi i të dhënave, që mund të jetë MySQL që nga versioni 5.7, PostgreSQL 9.6+, MongoDB 3.2+ ();
- klastri Apache Kafka;
- instancë Kafka Connect (versionet 1.x, 2.x);
- lidhësi Debezium i konfiguruar.
Punët për dy pikët e para, pra procesi i instalimit të DBMS dhe Apache Kafka, janë jashtë përmbajtjes së këtij artikulli. Megjithatë, për ata që dëshirojnë ta vendosin gjithçka në një ambient provues, në depo zyrtare me shembuj ka një .
Ne do të ndalemi më në detaje për dy pikët e fundit.
0. Kafka Connect
Këtu dhe më tej në këtë artikull, të gjitha shembujt e konfigurimit shqyrtohen në kontekstin e pamjen Docker, të shpërndarë nga zhvilluesit e Debezium. Ai përmban të gjitha skedaret e nevojshme të plugjineve (lidhësit) dhe parashikon konfigurimin e Kafka Connect nëpërmjet variablave të mjedisit.
Nëse pritet përdorimi i Kafka Connect nga Confluent, do të nevojitet që të shtoni manualisht plugjinat e lidhësve të nevojshëm në direktoriumin e caktuar në plugin.path ose të caktuar përmes variablës së mjedisit CLASSPATH. Konfigurimet e punëtorit Kafka Connect dhe lidhësve përcaktohen përmes skedarëve të konfigurimit, të cilat kalohen si argumente në komandën e nisjes së punëtorit. Më shumë detaje shihni në .
I gjithë procesi i konfigurimit të Debezium në variantin me lidhësin realizohet në dy etapa. Le të shqyrtojmë secilën prej tyre:
1. Konfigurimi i kornizës Kafka Connect
Për streamimin e të dhënave në klasterin Apache Kafka, në kornizën Kafka Connect vendosen parametra specifikë, të tilla si:
- parametrat e lidhjes me klasterin,
- emrat e topic-eve, në të cilat do të ruhet konfigurimi i vetë lidhësit,
- emri i grupit, në të cilin është nisur lidhësi (në rastin e përdorimit të mënyrës së shpërndarë).
Imazhi zyrtar Docker i projektit mbĂ«shtet konfigurimin pĂ«rmes variablave tĂ« mjedisit â dhe ne do ta pĂ«rdorim kĂ«tĂ«. Pra, shkarkoni imazhin:
docker pull debezium/connectGrupi minimal i variablave të mjedisit, i nevojshëm për të nisur lidhësin, duket si në vijim:
-
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â topic-i pĂ«r ruajtjen e pozicioneve tĂ« tĂ« cilat ndodhet aktualisht lidhĂ«si; -
CONNECT_STATUS_STORAGE_TOPIC=connector-statusâ topic pĂ«r ruajtjen e statusit tĂ« lidhĂ«sit dhe detyrave tĂ« tij; -
CONFIG_STORAGE_TOPIC=connector-configâ topic pĂ«r ruajtjen e tĂ« dhĂ«nave tĂ« konfigurimit tĂ« lidhĂ«sit dhe detyrave tĂ« tij; -
GROUP_ID=1â identifikuesi i grupit tĂ« punĂ«torĂ«ve ku mund tĂ« ekzekutohet detyra e lidhĂ«sit; i nevojshĂ«m kur pĂ«rdoret nĂ« modin e shpĂ«rndarĂ« (distributed) modi.
Po nisim 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 për Avro
Në mënyrë të paracaktuar, Debezium shkruan të dhënat në formatin JSON, që është i pranueshëm për sandboxing dhe volumet e vogla të të dhënave, por mund të bëhet problematike në baza të dhënash me ngarkesë të lartë. Një alternativë për konvertuesin JSON është serialization e mesazheve nëpërmjet formatit binar, i cili 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ë:
emri: CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL
vlera: http://kafka-registry-01:8081/
emri: CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL
vlera: http://kafka-registry-01:8081/
emri: VALUE_CONVERTER
vlera: io.confluent.connect.avro.AvroConverterDetajet mbi pĂ«rdorimin e Avro dhe konfigurimin e registrit janĂ« jashtĂ« pĂ«rmbajtjes sĂ« kĂ«tij artikulli â pĂ«r qĂ«llim ilustruar, 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Ă« shqyrtojmĂ« shembujt e konektorĂ«ve pĂ«r dy DB: PostgreSQL dhe MongoDB, pĂ«r tĂ« cilat kam pĂ«rvojĂ« dhe pĂ«r tĂ« cilat ka dallime (nĂ«se janĂ« tĂ« vogla, por nĂ« disa raste â thelbĂ«sore!).
Konfigurimi përshkruhet në notacionin JSON dhe ngarkohet në Kafka Connect përmes 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 kësaj konfigurimi është mjaft i thjeshtë:
- Në fillim, ai lidhet me bazën e të dhënave që është e caktuar në konfigurim dhe fillon në modalitetin snapshot inizial, duke dërguar në Kafka një sërë të dhënash fillestare, të marra përmes një
SELECT * FROM table_name. - Pas përfundimit të inicializimit, konektori kalon në modalitetin 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 (p.sh. pĂ«r tĂ« parĂ« statusin/ripĂ«rtrijĂ«/nĂ« pĂ«rditĂ«simin e konfigurimit) pĂ«rmes REST API tĂ« Kafka Connect; -
connector.classâ tipi i lidhĂ«sit tĂ« DBMS qĂ« do tĂ« pĂ«rdoret nga lidhĂ«si i konfiguruar; -
plugin.nameâ emri i plugin-it pĂ«r dekodimin logjik tĂ« tĂ« dhĂ«nave nga skedarĂ«t WAL. NĂ« dispozicion janĂ«wal2json,decoderbuffsdhepgoutput. Dy tĂ« parat kĂ«rkojnĂ« instalimin e pĂ«rkatĂ«sve nĂ« DBMS, ndĂ«rsapgoutputpĂ«r PostgreSQL versionin 10 dhe mĂ« lart nuk kĂ«rkon manipulime tĂ« tjera; -
database.*â opsionet pĂ«r lidhjen me DB, kudatabase.server.nameâ emri i instancĂ«s PostgreSQL, qĂ« pĂ«rdoret pĂ«r formimin e emrit tĂ« temĂ«s nĂ« klasterin Kafka; -
table.include.listâ lista e tabelave ku duam 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 lidhĂ«si dĂ«rgon mesazhe heartbeat nĂ« njĂ« temĂ« tĂ« veçantĂ«; -
heartbeat.action.queryâ pyetja qĂ« do tĂ« ekzekutohet me dĂ«rgesĂ«n e çdo mesazhi heartbeat (opsioni u shfaq me versionin 1.1); -
slot.nameâ emri i slotit tĂ« replikimit qĂ« do tĂ« pĂ«rdoret nga lidhĂ«si; publication.nameâ emri nĂ« PostgreSQL, qĂ« pĂ«rdor lidhĂ«si. NĂ« rast se ajo nuk ekziston, Debezium do tĂ« pĂ«rpiqet ta krijojĂ« atĂ«. NĂ«se pĂ«rdoruesi, nĂ«n tĂ« cilin po lidhet, nuk ka mjaft tĂ« drejta pĂ«r kĂ«tĂ« veprim, lidhĂ«si do tĂ« pĂ«rfundojĂ« punĂ«n me njĂ« gabim;-
transformonpërcakton se si do të ndryshojmë emrin e temës së synuar:-
transforms.AddPrefix.typetregon se do të përdorim shprehje të rregullta; -
transforms.AddPrefix.regexâ maska pĂ«r tĂ« cilĂ«n riemĂ«rohet emri i temĂ«s sĂ« synuar; -
transforms.AddPrefix.replacementâ nĂ« mĂ«nyrĂ« tĂ« drejtpĂ«rdrejtĂ« ajo, nĂ« tĂ« cilĂ«n riemĂ«rohet.
-
Më shumë mbi heartbeat dhe transformon
Në mënyrë të parazgjedhur, lidhësi dërgon të dhënat në Kafka për çdo transaksion të komituar, dhe numri i tij LSN (Log Sequence Number) shënon në temën ndihmëse offset. Por çfarë ndodh nëse lidhësi është i konfiguruar për të lexuar jo të gjithë bazën për tërësi, por vetëm një pjesë të tabelave të saj (ku azhurnimi i të dhënave ndodh jo shpesh)?
- Lidhësi do të lexojë skedarët WAL dhe nuk do të zbulojë në to komitin e transaksioneve në ato tabela, për të cilat po monitoron.
- Prandaj, ai nuk do të përditësojë pozitat 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 në ndoshta shterimin e tërë hapësirës diskore.
Dhe këtu vijnë në ndihmë opsionet heartbeat.interval.ms dhe heartbeat.action.query. Përdorimi i këtyre opsioneve në çift rezulton në një kërkesë për të ndryshuar të dhënat në një tabelë të veçantë çdo herë që dërgohet një mesazh pulse. Kështu, përditësohet vazhdimisht LSN, mbi të cilin ndodhet tani konektori (në slotin e replikimit). Kjo i lejon DBMS të fshijë skedarët WAL, që nuk janë më të nevojshëm. Më shumë mund të mësoni rreth funksionimit të opsioneve në .
Një opsion tjetër, që meriton më shumë vëmendje, është transformon. Megjithatë, ai është më shumë për lehtësinë dhe estetiken...
Prej default, Debezium krijon tema duke u bazuar në politikën e emërimit si vijon: serverName.schemaName.tableName. Kjo nuk është gjithmonë e përshtatshme. Me opsionet transformon mund të përcaktoni lista e tabelave, ngjarjet e të cilave duhet të rruhen në temën me një emër të caktuar, duke përdorur shprehje të rregullta.
NĂ« konfigurimin tonĂ«, falĂ« NĂ« pĂ«rfundim tĂ« pĂ«rshkrimit tĂ« konfigurimit tĂ« konnaktorit pĂ«r PostgreSQL, do tĂ« ishte e arsyeshme tĂ« flisnim pĂ«r karakteristikat/kufizimet e mĂ«poshtme tĂ« funksionimit tĂ« tij: Tani, do tĂ« ngarkojmĂ« konfigurimin tonĂ« nĂ« konnaktor: KontrollojmĂ« nĂ«se ngarkesa kaloi me sukses dhe konnaktori u aktivizua: Super: Ă«shtĂ« i konfiguruar dhe gati pĂ«r punĂ«. Tani do tĂ« sillen si konsumator dhe do tĂ« lidhim me Kafka, pastaj do tĂ« shtojmĂ« dhe ndryshojmĂ« njĂ« regjistrim nĂ« tabelĂ«: NĂ« temĂ«n tonĂ«, kjo do tĂ« shfaqet si mĂ« poshtĂ«: NjĂ« JSON shumĂ« i gjatĂ« me ndryshimet tona NĂ« tĂ« dyja rastet, regjistrimet pĂ«rbĂ«hen nga çelĂ«si (PK) i regjistrit qĂ« Ă«shtĂ« ndryshuar dhe vetĂ« natyra e ndryshimeve: si ishte regjistri para dhe si Ă«shtĂ« pas. Ky konektor pĂ«rdor mekanizmin standard tĂ« replikimit tĂ« MongoDB, duke lexuar informacionin nga oplog-i i nyjĂ«s kryesore tĂ« DBMS. Po ashtu si konektori i pĂ«rshkruar tashmĂ« pĂ«r PgSQL, kĂ«tu gjithashtu, gjatĂ« nisjes sĂ« parĂ«, merret njĂ« snapshot fillestar i tĂ« dhĂ«nave, pas sĂ« cilĂ«s konektori kalon nĂ« modin e leximit tĂ« oplog-it. Shembulli i konfiguracionit: Si mund tĂ« vĂ«rehet, kĂ«tu nuk ka opsione tĂ« reja krahasuar me shembullin e kaluar, por numri i opsioneve qĂ« pĂ«rfshijnĂ« lidhjen me DB-nĂ« dhe prefikset e tyre ka rĂ«nĂ«. CilĂ«simet Problemi i qĂ«ndrueshmĂ«risĂ« dhe disponibilitetit tĂ« lartĂ« Ă«shtĂ« mĂ« i rĂ«ndĂ«sishĂ«m se kurrĂ« â sidomos kur flasim pĂ«r tĂ« dhĂ«na dhe transaksione, dhe ndjekja e ndryshimeve tĂ« tĂ« dhĂ«nave nuk Ă«shtĂ« njĂ« çështje qĂ« e lĂ«mĂ« mĂ«njanĂ«. Le tĂ« shqyrtojmĂ« se çfarĂ« nĂ« parim mund tĂ« shkojĂ« gabim dhe çfarĂ« do tĂ« ndodhte me Debezium nĂ« çdo rast. Ka tre lloje dĂ«shtimesh: MegjithatĂ«, ka pĂ«rjashtime. NĂ«se konektori ka qenĂ« i çaktivizuar pĂ«r njĂ« periudhĂ« tĂ« gjatĂ« (ose nuk ka arritur tĂ« lidhĂ«t me instancĂ«n e MongoDB), dhe oplog gjatĂ« kĂ«saj kohe ka kaluar rotacionin, atĂ«herĂ« nĂ« rikthimin e lidhjes, konektori do tĂ« vazhdojĂ« tĂ« lexojĂ« tĂ« dhĂ«nat nga pozita e disponueshme e parĂ«, duke bĂ«rĂ« qĂ« njĂ« pjesĂ« e tĂ« dhĂ«nave nĂ« Kafka nuk tĂ« pĂ«rfshihet. Debezium â pĂ«rvoja ime e parĂ« me sistemet CDC dhe nĂ« tĂ«rĂ«si mjaft pozitive. Projekti mĂ« tĂ«rheq me mbĂ«shtetje pĂ«r DBMS-tĂ« kryesore, thjeshtĂ«sinĂ« e konfigurimit, mbĂ«shtetje pĂ«r klasterizimin dhe njĂ« komunitet aktiv. AtĂ«herĂ«, pĂ«r ata qĂ« janĂ« tĂ« interesuar pĂ«r praktikĂ«n, rekomandoj tĂ« shikoni udhĂ«zuesit pĂ«r dhe . Krahasuar me JDBC lidhĂ«sin pĂ«r Kafka Connect, pĂ«rparĂ«sia kryesore e Debezium Ă«shtĂ« se ndryshimet lexohen nga logjet e DBMS, duke lejuar marrjen e tĂ« dhĂ«nave me njĂ« vonesĂ« minimale. JDBC Connector (nga furnizimi i Kafka Connect) bĂ«n pyetje nĂ« tabelĂ«n e ndjekur me njĂ« interval tĂ« caktuar 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 zgjidhje tĂ« tjera (pĂ«rveç Debezium): Lexoni gjithashtu nĂ« blogun tonĂ«: Burimi: habr.comtransformon ndodhin kĂ«to: tĂ« gjitha ngjarjet CDC nga DB i ndjekur do tĂ« shkojnĂ« nĂ« temĂ«n me emrin data.cdc.dbname. NĂ« tĂ« kundĂ«rt, (pa kĂ«to cilĂ«sime) Debezium do tĂ« krijonte nĂ« mĂ«nyrĂ« standarde nga njĂ« topic pĂ«r çdo tabelĂ« si: pg-dev.public..
Kufizimet e konnaktorit
Zbatimi i konfiguracionit
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: vlera para (para) Ă«shtĂ« e barabartĂ« me null, ndĂ«rsa pas Ă«shtĂ« njĂ« string qĂ« Ă«shtĂ« futur. UPDATE: nĂ« payload.before tregohet gjendja e mĂ«parshme e stringut, ndĂ«rsa nĂ« payload.after â e re me natyrĂ«n 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 këtë herë bëjnë si më poshtë: kthejnë emrin e temës së synuar nga skema .. në data.cdc.mongo_.Qëndrueshmëria
Përfundimi
P.S.
