Hyrje nĂ« Debezium — CDC pĂ«r Apache Kafka

Hyrje nĂ« Debezium — CDC pĂ«r Apache Kafka

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?

Debezium — Ă«shtĂ« njĂ« pĂ«rfaqĂ«sues i kategorisĂ« sĂ« softuerit CDC (Capture Data Change), 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 Një projekt Open Source, 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:

Hyrje nĂ« Debezium — CDC pĂ«r Apache Kafka

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

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

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

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

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

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

Detajet 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, decoderbuffs dhe pgoutput. Dy tĂ« parat kĂ«rkojnĂ« instalimin e zgjatjeve pĂ«rkatĂ«se nĂ« DBMS, ndĂ«rsa pgoutput pĂ«r PostgreSQL versionin 10 e lart nuk kĂ«rkon veprime tĂ« tjera;
  • database.* — opsionet pĂ«r t'u lidhur me DB, ku database.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Ă« formatin schema.table_name; nuk mund tĂ« pĂ«rdoret sĂ« bashku me table.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 publikim 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;
  • transformon pĂ«rcakton se si tĂ« ndryshohet emri i temĂ«s sĂ« synuar:
    • transforms.AddPrefix.type tregon 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ë dokumentacionin.

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

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:

  1. Funksionaliteti i konektorit për PostgreSQL mbështetet në konceptin e dekodimit logjik. Prandaj, ai nuk ndjek kërkesat për ndryshimin e strukturës së BDs (DDL) - për rrjedhojë, në topicet e këtyre të dhënave nuk do të ketë.
  2. Duke pasur parasysh se përdoren slots replikimi, lidhja e konektorit është e mundur të me instancën kryesore të DBMS.
  3. Nëse përdoruesit, nën të cilin konektori lidhet me bazën e të dhënave, i janë dhënë të drejtat vetëm për lexim, atëherë para nisjes së parë është e nevojshme të krijoni manualisht një slot replikimi dhe publikimin në BD.

Përdorimi i konfigurimit

Pra, le të ngarkojmë konfigurimin tonë në konektor:

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

Kontrollojmë që ngarkimi ka kaluar me sukses dhe konektori ka nisur:

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

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

$ 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ë topicin tonë kjo do të paraqitet 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ë 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.

  • NĂ« rastin e SHTO: vlera para (para) Ă«shtĂ« e barabartĂ« me null, dhe pas — vargu qĂ« u fut.
  • NĂ« rastin e 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

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:

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

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

  1. Dështimi i Kafka Connect.. Nëse Connect është konfiguruar për të punuar në modin e shpërndarë, për këtë është e nevojshme të caktohen disa punëtorë me të njëjtin group.id. Atëherë, në rastin e dështimit të një prej tyre, konektori do të rinis në një punëtor tjetër dhe do të vazhdojë të lexojë nga pozita e fundit të komituar në temë në Kafka.
  2. Humbja e lidhjes me klasterin Kafka.. Konektori thjesht do të ndalojë leximin në pozitat që nuk arriti të dërgojë në Kafka, dhe do të përpiqet periodikisht të dërgojë atë përsëri, derisa përpjekja të përfundojë me sukses.
  3. Papërshtatshmëria e burimit të të dhënave. Konektori do të përpiqet të rikthejë lidhjen me burimin sipas konfigurimit. Në mënyrë të parazgjedhur, ky është 16 përpjekje duke përdorur vënien në prapavijë exponentiale. Pas përpjekjes së 16-të të dështuar, detyra do të shënohet si dështuar dhe do të kërkojë rinisje manuale përmes ndërfaqes REST të Kafka Connect.
    • NĂ« rastin e PostgreSQL tĂ« dhĂ«nat nuk do tĂ« humbasin, pasi pĂ«rdorimi i slot-eve tĂ« replikimit nuk do tĂ« lejojĂ« qĂ« tĂ« fshihen skedarĂ«t WAL, tĂ« papĂ«rpunuar nga konektori. NĂ« kĂ«tĂ« rast, ka edhe njĂ« anĂ« tjetĂ«r: nĂ«se lidhja rrjetike ndĂ«rmjet konektorit dhe DBMS-sĂ« ndĂ«rpritet pĂ«r njĂ« periudhĂ« tĂ« gjatĂ«, ka mundĂ«sinĂ« qĂ« hapsira nĂ« disk tĂ« mbarojĂ«, qĂ« mund tĂ« çojĂ« nĂ« dĂ«shtimin e plotĂ« tĂ« DBMS-sĂ«.
    • NĂ« rastin e MySQL skedarĂ«t e binlog-ut mund tĂ« rrotullohen nga vetĂ« DBMS mĂ« herĂ«t se sa tĂ« rikthehet lidhja. Kjo do tĂ« çonte nĂ« faktin qĂ« konektori do tĂ« kalojĂ« nĂ« gjendjen e dĂ«shtuar, dhe pĂ«r tĂ« rikthyer funksionimin normal do tĂ« kĂ«rkohet njĂ« rinisje nĂ« modin snapshot fillestar pĂ«r tĂ« vazhduar leximin nga binlog-Ă«t.
    • PĂ«r MongoDB. Dokumentacioni thotĂ«: sjellja e konektorit nĂ« rast se skedarĂ«t e regjistrove/oplog-ut janĂ« fshirĂ« dhe konektori nuk mund tĂ« vazhdojĂ« tĂ« lexojĂ« nga pikĂ«ku i ndaluar, Ă«shtĂ« e njĂ«jtĂ« pĂ«r tĂ« gjitha DBMS-tĂ«. Kjo nĂ«nkupton se konektori do tĂ« kalojĂ« nĂ« gjendjen dĂ«shtuar dhe do tĂ« kĂ«rkojĂ« rinisje nĂ« modin snapshot fillestar.

      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.

Përfundim

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

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

P.S.

Lexoni gjithashtu në blogun tonë:

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster