Njohja me Debezium — CDC pĂ«r Apache Kafka

Njohja me Debezium — CDC pĂ«r Apache Kafka

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?

Debezium — pĂ«rfaqĂ«sues i kategorisĂ« sĂ« softuerit CDC (Capture Data Change), dhe mĂ« saktĂ« — Ă«shtĂ« njĂ« koleksion konnektoresh pĂ«r DBMS tĂ« ndryshme, tĂ« pĂ«rshtatshme me kornizĂ«n Apache Kafka Connect.

Kjo Projekt Open Source, 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:

Njohja me Debezium — CDC pĂ«r Apache Kafka

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 embedded-connector. 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:

  1. burimi i të dhënave, që mund të jetë MySQL që nga versioni 5.7, PostgreSQL 9.6+, MongoDB 3.2+ (lista e plotë);
  2. klastri Apache Kafka;
  3. instancë Kafka Connect (versionet 1.x, 2.x);
  4. 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ë docker-compose.yaml.

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ë dokumentacion.

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/connect

Grupi 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.2

Shë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 Avro 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ë schema-registry (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.AvroConverter

Detajet 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, decoderbuffs dhe pgoutput. Dy tĂ« parat kĂ«rkojnĂ« instalimin e pĂ«rkatĂ«sve nĂ« DBMS, ndĂ«rsa pgoutput pĂ«r PostgreSQL versionin 10 dhe mĂ« lart nuk kĂ«rkon manipulime tĂ« tjera;
  • database.* — opsionet pĂ«r lidhjen me DB, ku database.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Ă« formatin schema.table_name; nuk mund tĂ« pĂ«rdoret sĂ« bashku me table.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 publikimin 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;
  • transformon pĂ«rcakton se si do tĂ« ndryshojmĂ« emrin e temĂ«s sĂ« synuar:
    • transforms.AddPrefix.type tregon 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ë dokumentacion.

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ë transformon 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

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:

  1. Funksionaliteti i konnaktorit për PostgreSQL mbështetet në konceptin e dekodimit logjik. Prandaj, ai nuk ndjek kërkesat për ndryshimin e strukturës së DB-së (DDL) - përkatësisht, në topicet e këtyre të dhënave nuk do të ketë.
  2. Duke qenë se përdoren slotet e replikimit, lidhja e konnaktorit është e mundur vetëm me instancën kryesore të DBMS-së.
  3. Nëse përdoruesit, me të cilin konnaktori lidhet me bazën e të dhënave, i janë dhënë të drejta vetëm për lexim, atëherë para nisjes së parë do të nevojitet krijimi manual i slotit të replikimit dhe publikimit në DB.

Zbatimi i konfiguracionit

Tani, do të ngarkojmë konfigurimin tonë në konnaktor:

curl -i -X POST -H "Accept:application/json" 
  -H  "Content-Type:application/json"  http://localhost:8083/connectors/ 
  -d @pg-con.json

Kontrollojmë nëse ngarkesa kaloi me sukses dhe konnaktori u aktivizua:

$ 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"}

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ë:

$ 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

Në temën tonë, kjo do të shfaqet si më poshtë:

Një JSON shumë i gjatë me ndryshimet tona

{
  "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
  }
}

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.

  • NĂ« rastin me INSERT: vlera para (para) Ă«shtĂ« e barabartĂ« me null, ndĂ«rsa pas Ă«shtĂ« njĂ« string qĂ« Ă«shtĂ« futur.
  • NĂ« rastin me 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

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:

{
  "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"
        }
  }

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 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

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:

  1. Dështimi i Kafka Connect. Nëse Connect është i konfiguruar për të punuar në një mënyrë të shpërndarë, nevojiten disa punonjës për t'i dhënë të njëjtin group.id. Kështu, kur njëri prej tyre dështon, lidhësi do të rinisë punën në një punonjës tjetër dhe do të vazhdojë të lexojë nga pozita e fundit të komituar në temën në Kafka.
  2. Humbja e lidhjes me klasterin Kafka. Lidhësi thjesht do të ndalë leximin në pozita që nuk arriti të dërgojë në Kafka dhe do të provojë të dërgojë përsëri me periudha të rregullta, deri sa përpjekja të jetë e suksesshme.
  3. Mungesës e burimit të të dhënave. Konnektori do të bëjë përpjekje për t'u rinisur me burimin sipas konfigurimit. Në mënyrë të parazgjedhur, kjo përfshin 16 përpjekje duke përdorur përplasje eksponenciale. Pas përpjekjes së 16-të të dështuara, detyra do të klasifikohet si e dështuar dhe do të kërkojë ri-nisje manuale përmes ndërfaqes REST të Kafka Connect.
    • NĂ« rastin me PostgreSQL tĂ« dhĂ«nat nuk do tĂ« humbasin, pasi pĂ«rdorimi i vendeve tĂ« replikuara nuk do tĂ« lejojĂ« qĂ« skedarĂ«t WAL tĂ« hiqen pa u lexuar nga konnektori. NĂ« kĂ«tĂ« rast ka edhe njĂ« anĂ« tjetĂ«r tĂ« medaljes: nĂ«se lidhja e rrjetit ndĂ«rmjet konnektorit dhe DBMS-sĂ« ndĂ«rpritet pĂ«r njĂ« periudhĂ« tĂ« gjatĂ«, ekziston rreziku qĂ« tĂ« pĂ«rfundojĂ« hapĂ«sira nĂ« disk, gjĂ« qĂ« mund tĂ« çojĂ« nĂ« dĂ«shtimin e plotĂ« tĂ« DBMS-sĂ«.
    • NĂ« rastin me MySQL skedarĂ«t binlog mund tĂ« rotullohen nga vetĂ« DBMS para se lidhja tĂ« rikthehet. Kjo do tĂ« çojĂ« nĂ« njĂ« gjendje dĂ«shtimi tĂ« konnektorit dhe pĂ«r tĂ« rikthyer funksionimin normal, do tĂ« nevojitet njĂ« ri-nisje nĂ« modin e snapshot fillestar pĂ«r tĂ« vazhduar leximin nga binlog.
    • PĂ«r MongoDB. Dokumentacioni thotĂ«: sjellja e konektorit nĂ« rast se skedarĂ«t e regjistrit/oplog u hoqĂ«n dhe konektori nuk mund tĂ« vazhdojĂ« leximin nga pozita ku ndaloi, Ă«shtĂ« e njĂ«jtĂ« pĂ«r tĂ« gjitha DBMS-tĂ«. Ajo konsiston nĂ« faktin se konektori do tĂ« kalojĂ« nĂ« gjendjen e dĂ«shtuar dhe do tĂ« kĂ«rkojĂ« ri-nisjen nĂ« modalitetin snapshot inizial.

      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.

Përfundimi

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 Kafka Connect dhe Debezium.

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):

P.S.

Lexoni gjithashtu në blogun tonë:

Burimi: habr.com

Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« đŸ”„ Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« | ProHoster