Tutvustus Debeziumisse — CDC Apache Kafka jaoks

Tutvustus Debeziumisse — CDC Apache Kafka jaoks

Oma töös kohtan ma tihti uusi tehnilisi lahendusi/programmide tooteid, mille kohta on vene keeles internetis suhteliselt vĂ€he teavet. Selle artikliga pĂŒĂŒan tĂ€ita ĂŒhe sellise lĂŒnga, tuues nĂ€ite oma hiljutisest praktikast, kus oli vajalik seadistada CDC-sĂŒndmuste saatmine kahest populaarsest andmebaasist (PostgreSQL ja MongoDB) Kafka klastrisse Debeziumi abil. Loodan, et see ĂŒlevaateartikkel, mis sĂŒndis tehtud töö tulemusena, osutub kasulikuks ka teistele.

Mis on Debezium ja ĂŒldiselt CDC?

Debezium — CDC tarkvara kategooria esindaja (Capture Data Change), tĂ€psemalt öeldes on see erinevate andmebaaside konnektorite kogum, mis on ĂŒhilduv Apache Kafka Connect raamistiku sĂŒsteemiga.

See Open Source projekt, millel on Apache License v2.0 litsents ja mida sponsoreerib ettevĂ”te Red Hat. Arendustööd on alustatud 2016. aastal ning praegu on ametlik tugi jĂ€rgmistele andmebaasidele: MySQL, PostgreSQL, MongoDB, SQL Server. Samuti on olemas konnektorid Cassandra ja Oracle jaoks, kuid need on praegu 'varajase juurdepÀÀsu' staatuses, ning uued vĂ€ljalasked ei garanteeri tagasipöördumatut ĂŒhilduvust.

Kui vĂ”rrelda CDC-d traditsioonilise lĂ€henemisega (kui rakendus loeb andmeid otse andmebaasist), siis selle peamisteks eelisteks on madala latentsusega, kĂ”rge usaldusvÀÀrsuse ja kĂ€ttesaadavuse taseme saavutamine andmete muutuste voogude realiseerimise kaudu ridade tasandil. Viimane kahest punktist saavutatakse, kasutades Kafka klastrit CDC-sĂŒndmuste salvestamiseks.

Samuti on eeliste hulka arvestatav, et sĂŒndmuste salvestamiseks kasutatakse ĂŒhtset mudelit, seega ei pea lĂ”pprakendus muretsema erinevate andmebaaside haldamise spetsiifikaga.

LÔpuks, thanks to the use of a message broker, the potential for horizontal scalability emerges for applications that track data changes. At the same time, the impact on the data source is minimized, as data is not retrieved directly from the database but rather from the Kafka cluster.

Debeziumi arhitektuurist

Debeziumi kasutamine piirdub sellise lihtsa skeemiga:

Andmebaas (andmeallikas) → konnektor Kafka Connectis → Apache Kafka → tarbija

Illustreerimiseks toome vÀlja skeemi projekti veebisaidilt:

Tutvustus Debeziumisse — CDC Apache Kafka jaoks

Kuid see skeem ei meeldi mulle eriti, kuna jÀÀb mulje, et on vÔimalik kasutada ainult sink-konnektorit.

Tegelikult on olukord teine: teie Data Lake'i tĂ€itmine (ĂŒlemises skeemis viimane element) — pole ainus viis Debeziumi rakendamiseks. Apache Kafka'sse saadetud sĂŒndmusi saab teie rakendustes kasutada erinevate olukordade lahendamiseks. NĂ€iteks:

  • mitteaktuaalsete andmete eemaldamine vahemĂ€lust;
  • teadete saatmine;
  • otsinguindeksite vĂ€rskendamine;
  • mingisugused auditilogid;
  • 


Kui teie rakendus on Java-s ja puudub vajadus/vÔimalus kasutada Kafka klastrit, on olemas ka vÔimalus töötada lÀbi embedded-konnektori. Selge eelis on see, et sellega ei pea lisainfrastruktuurist (konnektor ja Kafka) loobuma. Siiski on see lahendus 1.1 versioonist alates kuulutatud aegunuks (deprecated) ja selle kasutamist ei soovitata enam (tulevastes vÀljaannetes vÔib selle toe eemaldada).

KÀesolevas artiklis kÀsitletakse soovitatud arendajate arhitektuuri, mis tagab talitlushÀireteta toimimise ja skaleeritavuse.

Ühenduse konfiguratsioon

Et alustada muutuste jĂ€lgimist meie peamise varanduse – andmete – osas, on meil vaja:

  1. andmeallikat, milleks vÔivad olla MySQL alates versioonist 5.7, PostgreSQL 9.6+, MongoDB 3.2+ (tÀielik nimekiri);
  2. Apache Kafka klaster;
  3. Kafka Connecti instants (versioonid 1.x, 2.x);
  4. konfigureeritud Debeziumi ĂŒhendaja.

Kaks esimest punkti, st andmebaasi ja Apache Kafka installimisprotsess, jÀÀvad artikli huvist vÀljapoole. Siiski, neile, kes soovivad kÔik liivakasti kÀivitada, on ametlikus nÀidiste hoidlas olemas valmis docker-compose.yaml.

Keskendume aga kahe viimase punkti sisule.

0. Kafka Connect

Siin ja edaspidi artiklis kaalutakse kĂ”iki konfigureerimise nĂ€iteid Debeziumi arendajate levitatava Docker-pildi kontekstis. See sisaldab kĂ”iki vajalikke pistikfailide (ĂŒhendajate) faile ja vĂ”imaldab Kafka Connecti konfigureerimist keskkonnamuutujate abil.

Kui eeldatakse, et kasutatakse Confluenti Kafka Connecti, tuleb vajalike ĂŒhendajate pistikfailid iseseisvalt lisada kausta, mis on mÀÀratud plugin.path vĂ”i mÀÀratletud keskkonnamuutujaga CLASSPATH. Kafka Connecti töötlus ja konnektorite seadistused mÀÀratakse konfigureerimisfailide kaudu, mis antakse töötluse kĂ€ivitamise kĂ€su argumentidena. Rohkem teavet leiate dokumentatsioonis.

Kogu Debeizumi seadistamisprotsess konnektoriga toimub kahes etapis. Vaatame igaĂŒht neist:

1. Kafka Connect raamistiku seadistamine

Andmete voogedastamiseks Apache Kafka klastrisse Kafka Connect raamistiku kaudu mÀÀratakse spetsiifilised parameetrid, nagu:

  • ĂŒhenduse parameetrid klastriga,
  • teemade nimed, kus salvestatakse otse konnektori konfiguratsioon,
  • grupi nimi, milles konnektor töötab (juhul kui kasutatakse jaotatud reĆŸiimi).

Projektile ametlik Docker-pilt toetab konfigureerimist keskkonnamuutujate abil — seda me ka kasutame. Nii et laadime alla pildi:

docker pull debezium/connect

Minimum keskkonnamuutujate kogum, mis on vajalik konnektori kÀivitamiseks, on jÀrgmine:

  • BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092 — algne Kafka klastrite serverite loend, et saada tĂ€ispakkumine klastrite liikmetest;
  • OFFSET_STORAGE_TOPIC=connector-offsets — teema, mis salvestab hetkel konnektori asukohad;
  • CONNECT_STATUS_STORAGE_TOPIC=connector-status — konnektori ja tema ĂŒlesannete oleku ladustamise teema;
  • CONFIG_STORAGE_TOPIC=connector-config — konnektori ja tema ĂŒlesannete konfiguratsioonide andmete ladustamise teema;
  • GROUP_ID=1 — töötajate grupi identifikaator, kus konnektori ĂŒlesanne vĂ”ib toimuda; vajalik jaotatud (distributed) reĆŸiimi kasutamisel.

KĂ€ivitame konteineri nende muutujaitega:

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

MĂ€rkus Avro kohta

Vaikimisi kirjutab Debezium andmeid JSON-formaadis, mis on sobiv liivakastide ja vĂ€ikeste andmemahtude jaoks, kuid vĂ”ib osutuda probleemseks kĂ”rge koormusega andmebaasides. JSON-konverteri alternatiiviks on sĂ”numite serialiseerimine Avro binaarseks formaadiks, mis vĂ”imaldab vĂ€hendada I/O alamsĂŒsteemi koormust Apache Kafka-s.

Avro kasutamiseks on vajalik eraldi schema-registry (skeemide hoidmiseks). Konverteri muutujaid kasutatakse jÀrgmiselt:

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

Avro kasutamise ja registreerimise seadistamise ĂŒksikasjad ĂŒletavad artikli piire — edaspidi kasutame selguse huvides JSON-i.

2. Konnektori seadistamine

NĂŒĂŒd saame minna otse konnektori seadistamise juurde, mis loeb andmeid allikast.

Vaatame kahe andmebaasi, PostgreSQLi ja MongoDB, konnektoreid — nendes on mul kogemusi ja on teatud erinevusi (kuigi vĂ€ikesed, vĂ”ivad need mĂ”nel juhul olla mĂ€rkimisvÀÀrsed!).

Konfiguratsioon on kirjeldatud JSON-i notatsioonis ja laaditakse Kafka Connecti POST-requests abil.

2.1. PostgreSQL

PostgreSQL konnektori konfiguratsiooni nÀide:

{
  "name": "pg-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "plugin.name": "pgoutput",
    "database.hostname": "127.0.0.1",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "definitelynotpassword",
    "database.dbname" : "dbname",
    "database.server.name": "pg-dev",
    "table.include.list": "public.(.*)",
    "heartbeat.interval.ms": "5000",
    "slot.name": "dbname_debezium",
    "publication.name": "dbname_publication",
    "transforms": "AddPrefix",
    "transforms.AddPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.AddPrefix.regex": "pg-dev.public.(.*)",
    "transforms.AddPrefix.replacement": "data.cdc.dbname"
  }
}

Konnektori tööpĂ”himĂ”te pĂ€rast sellist seadistust on ĂŒsna lihtne:

  • Esimese kĂ€ivitamise ajal ĂŒhendub see konfigureeritud andmebaasiga ja siseneb reĆŸiimi esialgne snapshot, saates Kafka'sse esmased andmeĂŒksused, mis on saadud tingimuslikust SELECT * FROM table_name.
  • PĂ€rast initsialiseerimise lĂ”petamist liigub konnektor PostgreSQL WAL-failidest muudatuste lugemise reĆŸiimi.

Kasutatavate valikute kohta:

  • name — konnektori nimi, mille jaoks kasutatakse allpool kirjeldatud seadistust; tulevikus kasutatakse seda nime konnektoriga töötamiseks (nt oleku vaatamine/taaskĂ€ivitamine/seadistuse vĂ€rskendamine) Kafka Connect REST API kaudu;
  • connector.class — andmepood, mida konfigureeritav ĂŒhendus kasutab;
  • plugin.name — pistiku nimi, mis on mĂ”eldud andmete loogiliseks dekodeerimiseks WAL-failidest. Saadaval on jĂ€rgmised valikud: wal2json, decoderbuffs ja pgoutput. Esimene kaks vajavad vastava laienduse installimist andmebaasis, samas kui pgoutput PostgreSQL versiooni 10 ja uuemate jaoks ei vajata tĂ€iendavaid toiminguid;
  • database.* — ĂŒhenduse valikud andmebaasiga, kus database.server.name — PostgreSQL instantsi nimi, mida kasutatakse teema nime genereerimiseks Kafka klastris;
  • table.include.list — tabelite loetelu, kus soovime jĂ€lgida muudatusi; sisestatakse vormingus schema.table_name; ei saa kasutada koos table.exclude.list;
  • heartbeat.interval.ms — intervall (millisekundites), mil konnektor saadab heartbeat-sĂ”numeid spetsiaalsesse teema;
  • heartbeat.action.query — pĂ€ring, mida tehakse iga heartbeat-sĂ”numi saatmisel (valik, mis on saadaval alates versioonist 1.1);
  • slot.name — replikatsiooni sloti nimi, mida konnektor kasutab;
  • publication.name — nimi avalikustamist PostgreSQL-is, mida ĂŒhendus kasutab. Kui see puudub, proovib Debezium selle luua. Kui kasutajal, kelle all ĂŒhendus toimub, ei ole piisavalt Ă”igusi selle toimingu teostamiseks, lĂ”petab ĂŒhendus töö tĂ”rketeatega;
  • muundamised mÀÀrab, kuidas tĂ€pselt sihtteema nime muuta:
    • transforms.AddPrefix.type nĂ€itab, et kasutame regulaaravaldisi;
    • transforms.AddPrefix.regex — muster, mille jĂ€rgi muudame sihtteema nime;
    • transforms.AddPrefix.replacement — see, millele me selle muutume.

Rohkem teavet heartbeat'i ja muundamiste kohta

Vaikimisi saadab ĂŒhendus andmed Kafka'sse iga kinnitatud tehingu korral ning selle LSN (Log Sequence Number) salvestatakse teenuse teema offset. Aga mis juhtub, kui ĂŒhendus on seadistatud lugema mitte kogu andmebaasi, vaid ainult osa selle tabelitest (kus andmete vĂ€rskendamine ei toimu sageli)?

  • Ühendus hakkab lugema WAL-faile ega leia neist tehingute kinnitamist nendes tabelites, mille ĂŒle ta jĂ€lgib.
  • SeetĂ”ttu ei uuenda ta oma praegust positsiooni ei teemas ega replikatsiooni slotis.
  • See toimetab, et WAL-failid jÀÀvad kettale „kinni“ ja kogu kettaruumi ammendamine vĂ”ib juhtuda.

Siinkohal tulevad appi valikud heartbeat.interval.ms ja heartbeat.action.query. Nende valikute kasutamine koos vĂ”imaldab iga kord, kui saadetakse sĂŒdamepekslemise sĂ”num, teha pĂ€ring andmete muutmiseks eraldi tabelis. Sellisel viisil uuendatakse pidevalt LSN-i, kus konnektor praegu asub (replikatsiooni pesas). See vĂ”imaldab DBMS-il eemaldada WAL-failid, mis ei ole enam vajalikud. Lisainfot valikute töö kohta leiate dokumentatsioonis.

Teine valik, mis vÀÀrib suuremat tÀhelepanu, on muundamised. Kuigi see on pigem mugavuse ja ilu pÀrast...

Vaikimisi loob Debezium teemasid, jĂ€rgides jĂ€rgmist nimede mÀÀramise poliitikat: serverName.schemaName.tableName. See ei pruugi alati mugav olla. Valikute abil muundamised vĂ”ib regulaarselt vĂ€ljendit kasutades mÀÀrata tabelite nimekirja, millelt sĂŒndmused tuleb suunata konkreetse nimega teema.

Meie konfiguratsioonis tĂ€nu muundamised juhtuda jĂ€rgmine: kĂ”ik CDC-sĂŒndmused jĂ€lgitavast andmebaasist jĂ”uavad teema nimega data.cdc.dbnameKui need seaded pole paigas, siis Debezium loodi vaikimisi igale tabelile teema pĂ”hjal: pg-dev.public.

.

Konnektori piirangud

PostgreSQL konnektori konfiguratsiooni kirjelduse lÔpetamiseks tasub rÀÀkida jÀrgmistest omadustest/piirangutest selle toimimises:

  1. PostgreSQL konnektori funktsionaalsus toetub loogilise dekodeerimise kontseptsioonile. SeetĂ”ttu ei jĂ€lgi see andmebaasi struktuuri muutmise (DDL) pĂ€ringuid — seega ei sisaldu nendes andmete teemades. Kuna kasutatakse replikeerimissekke, on konnektori ĂŒhendamine vĂ”imalik
  2. peamise andmebaasi eksemplariga. ainult Kui kasutajal, kelle alt konnektor andmebaasi ĂŒhendust vĂ”tab, on antud ainult lugemisĂ”igused, siis tuleb enne esimest kĂ€ivitamist kĂ€sitsi luua replikeerimisvahekoht ja publikatsioon andmebaasis.
  3. Konfiguratsiooni rakendamine

Seega laadime meie konfigureerimise konnektorisse:

Kontrollime, et laadimine toimus edukalt ja konnektor kÀivitus:

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

Kontrollime, et laadimine Ônnestus ja konnektor kÀivitati:

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

SuurepĂ€rane: see on seadistatud ja valmis töötama. NĂŒĂŒd kĂ€itume tarbijana ja ĂŒhendame end Kafka'ga, seejĂ€rel lisame ja muudame tabelisse kirje:

$ 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

Meie teemas kuvab see jÀrgmiselt:

VĂ€ga pikk JSON meie muudatustega

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

MÔlemal juhul koosnevad kirjed jÀrjekorra (PK) vÔtmetest, mis on muudetud, ja muudatuste sisust: millisena oli kande sisu enne ja millisena ta pÀrast muutus.

  • Juhtumite puhul, INSERT: vÀÀrtus enne (before) on vĂ”rreldav null, ja pĂ€rast on rida, mis on lisatud.
  • Juhtumite puhul, KUUDA: payload.before kuvatakse rea eelmine seisund, samas kui payload.after on uus koos muudatuste sisuga.

2.2 MongoDB

See konnektor kasutab MongoDB standardset replikatsiooni mehhanismi, lugedes teavet pÔhivÔlvi oplog'ist.

Sarnaselt juba kirjeldatud PgSQL konnektorile, tehakse siin esmakordsel kĂ€ivitamisel andmete esmane snapshot, pĂ€rast mida vahetab konnektor ĂŒle oplog'i lugemise reĆŸiimile.

Konfiguratsiooni nÀidis:

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

Nagu mĂ€rkida, ei ole siin uusi vĂ”imalusi vĂ”rreldes eelmise nĂ€itega, kuid on vĂ€henenud vaid andmebaasiĂŒhenduste ja nende prefikside arv.

Seaded muundamised Seekord tehakse jĂ€rgmist: sihtteema nimi muudetakse skeemiks .. ĂŒhes data.cdc.mongo_.

Veadeta töö

TĂ”rke- ja kĂ”rge kĂ€ttesaadavuse kĂŒsimus on tĂ€napĂ€eval aktuaalne, eriti kui rÀÀgime andmetest ja tehingutest, ning andmete muudatuste jĂ€lgimine ei jÀÀ selles osas ĂŒle. Vaatame, mis vĂ”ib pĂ”himĂ”tteliselt valesti minna ning mis juhtub Debeziumiga igas neist olukordadest.

On kolm tÔrkevarianti:

  1. Kafka Connect'i tĂ”rge. Kui Connect on seadistatud töötama hajutatud reĆŸiimis, tuleb mitmele töötlejale mÀÀrata sama group.id. Sellisel juhul taaskĂ€ivitab ĂŒhendus ĂŒhe töötleja tĂ”rke korral end teisel töötlejal ning jĂ€tkab lugemist viimase kinnitatud positsiooniga Kafka teemas.
  2. Kaotatud ĂŒhendus Kafka-klastriga. Ühendus peatab lugemise positsioonilt, mida ei Ă”nnestunud saata Kafka'sse, ning proovib perioodiliselt uuesti saata, kuni katse on edukas.
  3. Andmete allika saadavuse puudumine. Ühendus pĂŒĂŒab andmeallikaga uuesti ĂŒhendust luua vastavalt seadistusele. Vaikimisi on see 16 katset, kasutades eksponentsiaalset taastekki. PĂ€rast 16. ebaĂ”nnestunud katset mĂ€rgitakse ĂŒlesanne kui ebaĂ”nnestunud ja selle kĂ€sitsi taaskĂ€ivitamine lĂ€bi REST-liidese Kafka Connect on vajalik.
    • Juhtumite puhul, PostgreSQL Andmed ei kao, kuna replikatsiooni slotide kasutamine ei luba WAL-faile kustutada, mida konnektor ei ole lugenud. Sellisel juhul on siiski ka negatiivne kĂŒlg: kui ĂŒhenduvus konnektori ja andmebaasi vahel katkeb pikaks ajaks, on oht, et kettaruumi saab otsa, mis vĂ”ib pĂ”hjustada andmebaasi tĂ€ieliku tĂ”rke.
    • Juhtumite puhul, MySQL Binaarfailide logid vĂ”ivad andmebaasi poolt enne ĂŒhenduvuse taastumist ringlusse vĂ”tta. See viib selleni, et konnektor lĂ€heb olekusse ebaĂ”nnestunud ja normaalse toimimise taastamiseks on vajalik uuesti kĂ€ivitamine algse pilti vĂ”ttes, et jĂ€tkata lugemist binaarfailidest.
    • Umbes MongoDB. Dokumentatsioon ĂŒtleb: konnektori kĂ€itumine juhul, kui logifailid / oplog on kustutatud ja konnektor ei saa jĂ€tkata lugemist positsioonilt, kus see peatunud oli, on kĂ”ikide DBMS-ide puhul sama. See seisneb selles, et konnektor lĂ€heb seisundisse ebaĂ”nnestunud ja nĂ”uab korduvkĂ€ivitamist reĆŸiimis esialgne snapshot.

      Kuid on erandeid. Kui konnektor on pikka aega olnud vĂ€ljas (vĂ”i ei saanud ĂŒhendust MongoDB eksemplariga), ja oplog on selle aja jooksul pöörlemise lĂ€binud, siis ĂŒhenduse taastamisel jĂ€tkab konnektor rahulikult andmete lugemist esimeselt kĂ€tte saadavalt positsioonilt, mistĂ”ttu osa andmeid jĂ”uab Kafka ei kuni.

KokkuvÔte

Debezium on minu esimene kogemus CDC-sĂŒsteemidega, ning ĂŒldiselt on see vĂ€ga positiivne. Projekt teeb mulje oma toetuse poolest peamistele DBMS-idele, konfigureerimise lihtsuse, klastritoe ja aktiivse kogukonna poolest. Praktikat huvitavatele soovitan tutvuda juhenditega Kafka Connect ja Debezium.

Debezium'i peamine eelis JDBC-ĂŒhendaja ees Kafka Connect'i jaoks on see, et muudatused loetakse andmebaasi logidest, mis vĂ”imaldab andmeid saada minimaalse viivitusega. JDBC Connector (Kafka Connect'i tarnekomplektist) teeb pĂ€ringuid jĂ€lgitava tabeli kohta kindla ajavahega ja (sama pĂ”hjusel) ei genereeri see sĂ”numeid andmete kustutamisel (kuidas saab kĂŒsida andmeid, mida enam pole?).

Sarnaste probleemide lahendamiseks tasub vaadata jÀrgmisi lahendusi (Lisaks Debeziumile):

P.S.

Lugege ka meie blogist:

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster