Tutvustame Debeziumi — CDC Apache Kafka jaoks

Tutvustame Debeziumi — CDC Apache Kafka jaoks

Oma töös kohtan sageli uusi tehnilisi lahendusi/softtooteid, mille kohta on venekeelses internetis ĂŒsna vĂ€he teavet. Selle artikliga pĂŒĂŒan tĂ€ita ĂŒht sellist puudujÀÀki, tuues nĂ€ite oma hiljutisest praktikast, kus tuli seadistada CDC-sĂŒndmuste saatmine kahest populaarset andmebaasist (PostgreSQL ja MongoDB) Kafka klastrisse Debeziumi abil. Loodan, et see ĂŒlevaade, mis on kirjutatud tehtud töö tulemuste pĂ”hjal, osutub kasulikuks ka teistele.

Mis on Debezium ja ĂŒldiselt CDC?

Debezium on CDC (Capture Data Change) tarkvara kategooria esindaja, tĂ€psemalt on see konnektorite kogum erinevatele andmebaasidele, mis on ĂŒhilduvad Apache Kafka Connect raamistikuga.

See on Avaallika projekt, millel on Apache License v2.0 litsents ja mille sponsoreerib Red Hat. Arendustööd on alustatud 2016. aastal ja praeguseks on selles ametlik tugi jĂ€rgmistele andmebaasidele: MySQL, PostgreSQL, MongoDB, SQL Server. Samuti on olemas konnektorid Cassandra ja Oracle jaoks, kuid praegu on need "varajase juurdepÀÀsu" staatuses ning uued vĂ€ljaanded ei taga tagasipöörduvat ĂŒhilduvust.

CDC-d vĂ”rreldes traditsioonilise lĂ€henemisega (kui rakendus loeb andmeid otse andmebaasist) on selle peamisteks eelisteks andmete muutuste streamimise realiseerimine ridade tasemel koos madala latentsuse, kĂ”rge usaldusvÀÀrsuse ja kĂ€ttesaadavusega. Viimane kaks punkti saavutatakse, kasutades CDC-sĂŒndmuste hoidmiseks Kafka klastrit.

Sarnaselt on eelisteks ka see, et sĂŒndmuste salvestamiseks kasutatakse ĂŒhtset mudelit, mistĂ”ttu lĂ”pp-rakendusel ei pea olema mures erinevate andmebaaside kasutamise nĂŒansside ĂŒle.

L finalmente, tÀnu sÔnumite vahendaja kasutamisele avardub rakenduste horisontaalne skaleerimine, mis jÀlgivad andmete muutusi. Samuti on andmeallika mÔju minimaalselt piiratud, kuna andmeid saadakse mitte otse andmebaasist, vaid Kafka klastrist.

Debeziumi arhitektuurist

Debeziumi kasutamine seisneb sellises lihtsas skeemis:

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

Illustreerimiseks toome vÀlja skeemi projekti veebilehelt:

Tutvustame Debeziumi — CDC Apache Kafka jaoks

Kuid see skeem mulle ei meeldi, kuna mul on tunne, et on vÔimalik ainult sink-konektori kasutamine.

Kuid tegelikult on olukord teistsugune: teie Data Lake'i tĂ€itmine (viimane element ĂŒlaltoodud skeemis) — see ei ole ainus viis Debeziumi rakendamiseks. Apache Kafkasse edastatud sĂŒndmusi saavad teie rakendused kasutada erinevate olukordade lahendamiseks. NĂ€iteks:

  • vĂ€heharjutatud andmete kustutamine vahemĂ€lust;
  • teavituste saatmine;
  • otsingumootorite indeksite vĂ€rskendamine;
  • mingi sarnane auditi logi;
  • 


Kui teil on Java-rakendus ega ole vaja/vÔimalust kasutada Kafka klastrit, on samuti vÔimalik kasutada embedded-konektorit. Ilmselge pluss on see, et selle abil saab vÀltida lisainfrastruktuuri (nagu konektor ja Kafka) kasutamist. Kuid see lahendus on versioonist 1.1 deprekeeritud ja rohkem ei soovitata. Tulevastes vÀljaannetes vÔib selle toe eemaldada.

Selles artiklis kÀsitletakse soovitatavat arhitektuuri, mille on vÀlja töötanud arendajad, et tagada tÔrkeoleku kindlus ja skaleeritavus.

Konektori konfiguratsioon

Andmete, millel on kÔige suurem vÀÀrtus, muudatuste jÀlgimiseks, on meil vajalik:

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

Esimese kahe punkti, st andmebaasi ja Apache Kafkade installimise protsess, ĂŒletab artikli piirid. Kuid neile, kes soovivad kĂ”ike liivakastis kĂ€ivitada, on ametlikus nĂ€idiste hoidlas valmis docker-compose.yaml.

Keskendume aga kahele viimasele punktile.

0. Kafka Connect

Siin ja edaspidi artikkel kÔik konfiguratsiooni nÀited kÀsitletakse Debeziumi arendajate levitatud Docker-pildis. See sisaldab kÔiki vajalikke pluginifailisid (konektorid) ning eeldab Kafka Connect'i seadistamist keskkonnamuutujate kaudu.

Kui kavatsete kasutada Confluenti Kafka Connect'i, peate ise vajalikud konektorite pluginad lisama katalooge, nagu on nÀidatud plugin.path vÔi seadistada keskkonnamuutujaga CLASSPATH. Kafka Connecti ja konnektorite seaded mÀÀratakse konfiguratsioonifailide kaudu, mis edastatakse argumendina töötaja kÀivitamise kÀsule. Lisateavet vt dokumentatsioon.

Kogu Debeziumi seadistamise protsess konnektori variandiga toimub kahes etapis. Vaatame igaĂŒht neist:

1. Kafka Connecti raamistikku seadistamine

Andmete voogedastamiseks Apache Kafka klastrisse seadistatakse Kafka Connecti raamistikus spetsiifilised parameetrid, nagu:

  • ĂŒhenduse parameetrid klastri juurde,
  • teemade nimed, kus konnektori seadistused endid hoitakse,
  • grupi nimi, milles konnektor on kĂ€ivitatud (jaotatud reĆŸiimi kasutamisel).

Projekti ametlik Docker-pilt toetab seadistamist keskkonna muutujate abil — seda me ka kasutame. Nii et, tĂ”mmake pilt:

docker pull debezium/connect

Minimaalne muutujate komplekt, mis on vajalik konnektori kÀivitamiseks, nÀeb vÀlja jÀrgmine:

  • BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092 — klastrisse Kafka ĂŒhendusserversite esialgne loend, et saada tĂ€ielik loend klastriliikmetest;
  • OFFSET_STORAGE_TOPIC=connector-offsets — teema, kus hoitakse positsioone, kus konnektor hetkel asub;
  • CONNECT_STATUS_STORAGE_TOPIC=connector-status — teema, kus hoitakse konnektori ja selle ĂŒlesannete staatust;
  • CONFIG_STORAGE_TOPIC=connector-config — teema, kus hoitakse konnektori ja selle ĂŒlesannete seadistuste andmeid;
  • GROUP_ID=1 — töötajate grupi ID, millel konnektori ĂŒlesanne vĂ”ib töötada; vajalik jaotatud (distributed) reĆŸiimi kasutamisel.

KĂ€ivita konteiner nende muutujatega:

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 andmed JSON-formaadis, mis sobib liivakastide ja vĂ€ikeste andmemahtude jaoks, kuid vĂ”ib muutuda probleemiks kĂ”rgelt koormatud andmebaasides. Alternatiiv JSON-konverterile on sĂ”numite serialiseerimine Avro binaarformaati, mis vĂ”imaldab vĂ€hendada I/O alamsĂŒsteemi koormust Apache Kafka-s.

Avro kasutamiseks on vajalik eraldi schema-registry (skeemide hoidmiseks). Muutujad konverteri jaoks nÀevad vÀlja jÀrgmised:

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 registri seadistamise ĂŒksikasjad jÀÀvad artikli raames vĂ€ljapoole — edaspidi kasutame selguse huvides JSON-i.

2. Konnektori seadistamine

NĂŒĂŒd saame liikuda otse konnektori konfiguratsiooni juurde, mis loeb andmeid allikast.

Vaatame kahe andmebaasi konnektori nĂ€itel: PostgreSQL ja MongoDB — mul on nende kohta kogemus ning nendes on erinevusi (kuigi vĂ€ikseid, vĂ”ivad need mĂ”nel juhul olla olulised!).

Konfiguratsioon on kirja pandud JSON-notatsioonis ja laaditakse Kafka Connecti POST-pÀringu kaudu.

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 seadistamist on ĂŒsna lihtne:

  • Esimese kĂ€ivitamise korral ĂŒhendub see konfiguratsioonis mĂ€rgitud andmebaasi ja kĂ€ivitub reĆŸiimis initial snapshot, saates Kafka-sse algse andmekogumi, mis on saadud tingimusliku SELECT * FROM table_name.
  • PĂ€rast initsialiseerimise lĂ”ppemist lĂ€heb konnektor ĂŒle muutuste lugemise reĆŸiimile PostgreSQL WAL-failidest.

Kasutatavatest valikutest:

  • nimi — konnektori nimi, mille jaoks kasutatakse allpool kirjeldatud konfiguratsiooni; hiljem kasutatakse seda nime konnektoriga töötamiseks (st staatuse vaatamiseks/kĂ€ivitamiseks/konfiguratsiooni vĂ€rskendamiseks) lĂ€bi Kafka Connect REST API;
  • connector.class — andmebaasi konnektori klass, mida konfigureeritav konnektor kasutab;
  • plugin.name — pluginanimi, mis on vajalik andmete loogiliseks dekodeerimiseks WAL-failidest. Saadaval on valikud wal2json, decoderbuffs ja pgoutput. Esimesed kaks nĂ”uavad vastavate laienduste installimist andmebaasi, kuid pgoutput PostgreSQL versioon 10 ja uuem ei vaja tĂ€iendavaid toiminguid;
  • database.* — valikud andmebaasi ĂŒhendamiseks, kus database.server.name — PostgreSQL instantsi nimi, mida kasutatakse teema nime loomiseks Kafka klastris;
  • table.include.list — nimekiri tabelitest, kus soovime muudatusi jĂ€lgida; mÀÀratakse formaadis schema.table_name; ei tohi kasutada koos table.exclude.list;
  • heartbeat.interval.ms — intervall (millisekundites), mil konnektor saadab heartbeat-sĂ”numeid spetsiaalsesse teema;
  • heartbeat.action.query — pĂ€ring, mis tĂ€idetakse iga heartbeat-sĂ”numi saatmise ajal (valik ilmus versioonis 1.1);
  • slot.name — replikatsioonislot'i nimi, mida konnektor kasutab;
  • publication.name — nimi vĂ€ljaandmise PostgreSQL-is, mida konnektor kasutab. Kui seda ei eksisteeri, proovib Debezium selle luua. Kui kasutajal, kelle alt ĂŒhendus toimub, ei ole selle toimingu jaoks piisavalt Ă”igusi, lĂ”petab konnektor töö veaga;
  • transforms mÀÀrab, kuidas tĂ€pselt muuta sihteema nime:
    • transforms.AddPrefix.type nĂ€itab, et kasutame regulaaravaldusi;
    • transforms.AddPrefix.regex — mask, mille alusel sihteema nimi mÀÀratakse uuesti;
    • transforms.AddPrefix.replacement — see, millele me selle ĂŒmber mÀÀrame.

Rohkem infot heartbeat'i ja transforms'i kohta

Vaikimisi saadab konnektor andmeid Kafka-sse igas komitendis tehingus ning selle LSN (Log Sequence Number) salvestatakse teenindusteemasse offset. Kuid mis juhtub, kui konnektor on seadistatud lugema mitte kogu andmebaasi, vaid ainult osa selle tabelitest (kus andmete vÀrskendamine ei toimu sageli)?

  • Konnektor loeb WAL-faile ja ei leia neist tehingute komiteid nendesse tabelitesse, millele ta tĂ€helepanu pöörab.
  • SeetĂ”ttu ei uuenda ta oma praegust positsiooni ei teemadesse ega replikatsioonislotis.
  • See omakorda toob kaasa WAL-failide "hoidmise" kettal ja vĂ”imaliku kogu kettaruumi ammendumise.

Ja siin tulevad appi valikud heartbeat.interval.ms ja heartbeat.action.query. Nende valikute kooskasutamine vÔimaldab igal ajal heartbeat-sÔnumi saatmisel teostada pÀringu andmete muutmiseks eraldi tabelis. Sellega ajakohastatakse pidevalt LSN-i, kus konnektor praegu asub (replikatsioonislotis). See vÔimaldab andmebaasi halduril kustutada WAL-failid, mis enam ei ole vajalikud. Lisainfot nende valikute töötamisest saab dokumentatsioon.

Teine vĂ”imalus, mis vÀÀrib suuremat tĂ€helepanu, on transforms. Kuigi see on pigem mugavuse ja ilu kohta


Debezium loob vaikimisi teemasid, jĂ€rgides jĂ€rgmisi nimetamisreegleid: serverName.schemaName.tableName. See ei pruugi alati mugav olla. VĂ”imalustega transforms saab kasutada regulaaravaldisi, et mÀÀrata tabelite loend, mille sĂŒndmusi peab suunama kindla nimega teema.

Meie konfiguratsioonis toimub tĂ€nu transforms jĂ€rgnev: kĂ”ik CDC-sĂŒndmused jĂ€lgitavast andmebaasist jĂ”uavad teema nimega data.cdc.dbname. Vastupidisel juhul (ilma nende seadistusteta) oleks Debezium vaikimisi loonud iga tabeli jaoks teema, mille formaat on: pg-dev.public.

.

Ühenduse piirangud

PostgreSQL konnektori konfiguratsiooni kokkuvÔtteks tasub rÀÀkida jÀrgmistest omadustest/piirangutest:

  1. PostgreSQL konnektori funktsionaalsus pĂ”hineb loogilise dekodeerimise kontseptsioonil. Seega ei jĂ€lgi ta andmebaasi struktuuri muutmise pĂ€ringuid (DDL) — seetĂ”ttu ei ole nende andmete teemades.
  2. Kuna kasutatakse replikatsiooniplokke, on konnektori ĂŒhendamine vĂ”imalik seda peamise andmebaasi instantsiga.
  3. Kui andmebaasi kasutajal, kelle all konnektor andmebaasi ĂŒhendub, on ainult lugemisĂ”igused, tuleb enne esimest kĂ€ivitamist kĂ€sitsi luua replikatsiooniplokk ja vĂ€ljaanne andmebaasis.

Konfiguratsiooni rakendamine

Seega laadime meie konfiguratsiooni konnektorisse:

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

Kontrollime, et laadimine lÀks hÀsti 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 tööks valmis. NĂŒĂŒd teeme nĂ€gu, et oleme tarbija ja ĂŒhendume Kafka'ga, seejĂ€rel lisame ja muudame kirjet tabelis:

$ 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 peegelduvad need jÀrgmise kujul:

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 salvestused kirje (PK) vÔtmetest, mis on muudetud, ja muudatuste olemusest: milline oli kirje enne ja milline pÀrast.

  • Juhtumi puhul INSERT: vÀÀrtus enne (before) on vĂ”rdne null, ja pĂ€rast on sisestatud string.
  • Juhtumi puhul UPDATE: payload.before kuvatakse rea eelmine olek, ja payload.after on uus koos muudatuste olemusega.

2.2 MongoDB

See konnektor kasutab MongoDB standardset replikatsiooni mehhanismi, lugedes teavet primaarse andmebaasi oplog’ist.

Sarnaselt eelnevalt kirjeldatud PgSQL konnektorile, teeb ka see esmakordsel kĂ€ivitamisel esialgse andmete snapshotsi, pĂ€rast mida lĂŒlitub konnektor 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Ă€rgata, ei ole siin uusi vĂ”imalusi vĂ”rreldes eelmise nĂ€itega, kuid on vĂ€henenud vaid andmebaasi ĂŒhendamise ja nende eeltehtud nimetuste arv.

Seaded transforms seekord teevad nad jÀrgmist: muudavad sihttopiku nime mustrist .. ja data.cdc.mongo_.

Talitluskatkestustunne

Vastus vastupidavuse ja kĂ”rge k disponibiliteedi kĂŒsimusele on tĂ€napĂ€eval ÀÀrmiselt oluline — eriti kui rÀÀgime andmetest ja tehingutest, ning andmete muutuste jĂ€lgimine ei jÀÀ selles kĂŒsimuses kĂ”rvale. Vaatleme, mis vĂ”ib tegelikult valesti minna ja mis toimub Debeziumiga igas juhtumis.

On kolm tÔrkevarianti:

  1. Kafka Connecti tĂ”rge. Kui Connect on seadistatud töötama hajutatud reĆŸiimis, peavad mitmed töötajad omama sama group.id. Sel juhul taastatakse konnektor teisel töötajal ja jĂ€tkab lugemist viimase kinnitatud asukohast Kafka topikus.
  2. Ühenduse kadumine Kafka klastriga. Koonnektor lihtsalt peatab lugemise asukohas, mida ei Ă”nnestunud saata Kafka'sse, ja ĂŒritab perioodiliselt uuesti saata, kuni katse Ă”nnestub.
  3. Andmeallika kĂ€ttesaamatust. Ühendus pĂŒĂŒab uuesti ĂŒhendada allika kĂŒlge vastavalt seadistusele. Vaikimisi on see 16 katset, kasutades eksponentsiaalset tagasihoidlikkust. PĂ€rast 16. ebaĂ”nnestunud katset mĂ€rgitakse ĂŒlesanne ebaĂ”nnestunuks ja selle taastamiseks on vajalik kĂ€sitsi kĂ€ivitamine lĂ€bi Kafka Connect REST-API.
    • Juhtumi puhul PostgreSQL Andmed ei kao, kuna replikatsiooni sildid ei luba WAL-faile eemaldada, mida ĂŒhendus ei ole lugenud. Sellega kaasneb aga ka varjukĂŒlg: kui ĂŒhendus ĂŒhenduse ja andmebaasi vahel katkeb pikemaks ajaks, on tĂ”enĂ€olisem, et ketasruum saab otsa, mis vĂ”ib viia andmebaasi tĂ€ieliku tĂ”rkeni.
    • Juhtumi puhul MySQL Binaarsete logifailide vĂ”ib andmebaas ise enne taastumise ĂŒhendust pöörata. See toob kaasa, et ĂŒhendus lĂ€heb ebaĂ”nnestumiseks ja normaalse töö taastamiseks on vaja uuesti kĂ€ivitada algse jÀÀgi reĆŸiimis, et jĂ€tkata lugemist binaarsetest logidest.
    • KĂŒsimus MongoDB. Dokumentatsioon ĂŒtleb: ĂŒhenduse kĂ€itumine juhul, kui logifailid/oplog on kustutatud ja ĂŒhendus ei saa lugemist jĂ€tkata seal, kus see jĂ€i, on kĂ”igi andmebaaside jaoks sama. See tĂ€hendab, et ĂŒhendus lĂ€heb olekusse ebaĂ”nnestunuks ja nĂ”uab taaskĂ€ivitamist reĆŸiimis initial snapshot.

      Kuid on ka erandeid. Kui ĂŒhendus on pikka aega olnud vĂ€lja lĂŒlitatud (vĂ”i ei ole saanud MongoDB instantsile ligi), ja oplog on selle aja jooksul pööratud, siis ĂŒhenduse taastumisel jĂ€tkab ĂŒhendus vaikselt andmete lugemist esimesest kĂ€ttesaadavast kohast, mille tĂ”ttu osa andmeid jĂ”uab Kafka. ei kaotsi.

KokkuvÔte

Debezium on minu esimene kogemus CDC-sĂŒsteemidega ja ĂŒldiselt vĂ€ga positiivne. Projekt paelus mind peamiste andmebaaside toega, lihtsa seadistamise, klastritoe ja aktiivse kogukonnaga. Huvi korral soovitan tutvuda juhenditega Kafka Connect ja Debezium.

Debeziumi peamine eelis Kafka Connect JDBC-ĂŒhendaja ees on see, et muudatused loetakse vĂ€lja andmebaasi pĂ€evikutest, mis vĂ”imaldab andmeid saada minimaalse viivitusega. JDBC-ĂŒhendaja (mis on osa Kafka Connectist) pĂ€rib jĂ€lgitavast tabelist kindla aja jooksul ja (sama pĂ”hjuse tĂ”ttu) ei genereeri teateid andmete kustutamisel (kuidas saab kĂŒsida andmeid, mida pole?).

Sarnaste probleemide lahendamiseks tasub vaadata jÀrgmisi lahendusi (vÀlja arvatud Debezium):

P.S.

Lugege ka meie blogist:

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster