Wprowadzenie do Debezium — CDC dla Apache Kafka

Wprowadzenie do Debezium — CDC dla Apache Kafka

W swojej pracy często napotykam na nowe rozwiązania techniczne lub produkty oprogramowania, o których w polskojęzycznym internecie jest stosunkowo mało informacji. W tym artykule postaram się uzupełnić jedną z tych luk, przywołując przykład z mojej niedawnej praktyki, kiedy to konieczne było skonfigurowanie wysyłania zdarzeń CDC z dwóch popularnych baz danych (PostgreSQL i MongoDB) do klastra Kafka za pomocą Debezium. Mam nadzieję, że ten przeglądowy artykuł, powstały na podstawie przeprowadzonej pracy, będzie również przydatny innym.

Czym jest Debezium i w ogóle CDC?

Debezium to przedstawiciel kategorii oprogramowania CDC (Capture Data Change), a dokładniej — to zestaw konektorów do różnych baz danych, które są zgodne z frameworkiem Apache Kafka Connect.

To Projekt Open Source, który korzysta z licencji Apache License v2.0 i jest sponsorowany przez firmę Red Hat. Rozwój trwa od 2016 roku i obecnie oferuje oficjalne wsparcie dla następujących baz danych: MySQL, PostgreSQL, MongoDB, SQL Server. Istnieją również konektory dla Cassandry i Oracle, ale w tej chwili znajdują się w stanie „wczesnego dostępu”, a nowe wydania nie gwarantują wstecznej kompatybilności.

Porównując CDC z tradycyjnym podejściem (gdy aplikacja bezpośrednio odczytuje dane z bazy danych), jego głównymi zaletami są realizacja strumieniowania zmian danych na poziomie wierszy z niskim opóźnieniem oraz wysoką niezawodnością i dostępnością. Te dwa ostatnie punkty osiąga się dzięki wykorzystaniu klastra Kafka jako magazynu zdarzeń CDC.

Inną zaletą jest to, że do przechowywania zdarzeń używana jest jednolita model, więc końcowa aplikacja nie musi martwić się o niuanse eksploatacji różnych baz danych.

Wreszcie, dzięki zastosowaniu brokera wiadomości otwiera się możliwość poziomego skalowania aplikacji, które monitorują zmiany danych. Przy tym wpływ na źródło danych jest minimalny, ponieważ pobieranie danych odbywa się nie bezpośrednio z bazy danych, ale z klastra Kafka.

O architekturze Debezium

Użycie Debezium sprowadza się do takiej prostej schemy:

Baza danych (jako źródło danych) → konektor w Kafka Connect → Apache Kafka → konsument

Jako ilustrację przedstawiam schemat z witryny projektu:

Wprowadzenie do Debezium — CDC dla Apache Kafka

Jednak ten schemat niezbyt mi się podoba, ponieważ stwarza wrażenie, że można używać tylko sink-connectora.

W rzeczywistości jednak sytuacja jest inna: zapełnienie waszego Data Lake (ostatni element na powyższym schemacie) — to nie jedyny sposób wykorzystania Debezium. Zdarzenia wysyłane do Apache Kafka mogą być wykorzystywane przez twoje aplikacje do rozwiązania różnych sytuacji. Na przykład:

  • usuwanie nieaktualnych danych z cache'a;
  • wysyłanie powiadomień;
  • aktualizacje indeksów wyszukiwania;
  • pewny rodzaj logów audytu;

W przypadku, gdy masz aplikację napisaną w Java i nie ma potrzeby/możliwości korzystania z klastra Kafka, istnieje również możliwość pracy przez embedded-connector. Oczywistą zaletą jest to, że można zrezygnować z dodatkowej infrastruktury (w postaci connectora i Kafka). Jednak to rozwiązanie zostało uznane za przestarzałe (deprecated) w wersji 1.1 i nie jest już zalecane do użycia (w przyszłych wydaniach wsparcie może zostać usunięte).

W tym artykule omawiana będzie zalecana przez deweloperów architektura, która zapewnia odporność na błędy i możliwość skalowania.

Konfiguracja connectora

Aby rozpocząć śledzenie zmian najważniejszego zasobu — danych, potrzebujemy:

  1. źródła danych, którym może być MySQL od wersji 5.7, PostgreSQL 9.6+, MongoDB 3.2+ (pełna lista);
  2. klaster Apache Kafka;
  3. instancję Kafka Connect (wersje 1.x, 2.x);
  4. odpowiednio skonfigurowany connector Debezium.

Prace nad dwoma pierwszymi punktami, tj. proces instalacji DBMS i Apache Kafka, wykraczają poza zakres tego artykułu. Jednak dla tych, którzy chcą uruchomić wszystko w piaskownicy, w oficjalnym repozytorium z przykładami znajduje się gotowy docker-compose.yaml.

Zatrzymamy się bardziej szczegółowo na dwóch ostatnich punktach.

0. Kafka Connect

W tym i w dalszej części artykułu wszystkie przykłady konfiguracji są rozważane w kontekście obrazu Docker, rozpowszechnianego przez deweloperów Debezium. Zawiera on wszystkie niezbędne pliki wtyczek (connectorów) i przewiduje konfigurację Kafka Connect za pomocą zmiennych środowiskowych.

W przypadku, gdy planowane jest wykorzystanie Kafka Connect od Confluent, konieczne będzie samodzielne dodanie wtyczek niezbędnych connectorów do katalogu podanego w plugin.path lub określanego poprzez zmienną środowiskową CLASSPATH. Ustawienia workera Kafka Connect i konektorów określane są za pomocą plików konfiguracyjnych, które są przekazywane jako argumenty podczas uruchamiania workera. Szczegóły w dokumentacji.

Cały proces konfiguracji Debezium z wykorzystaniem konektora odbywa się w dwóch etapach. Przeanalizujmy każdy z nich:

1. Konfiguracja frameworka Kafka Connect

Aby przesyłać dane do klastra Apache Kafka w frameworku Kafka Connect, definiuje się specyficzne parametry, takie jak:

  • parametry połączenia z klastrem,
  • nazwy topiców, w których będzie przechowywana sama konfiguracja konektora,
  • nazwa grupy, w której uruchomiony jest konektor (w przypadku użycia trybu rozproszonego).

Oficjalny obraz Docker projektu wspiera konfigurację za pomocą zmiennych środowiskowych — z tego skorzystamy. Zatem, pobieramy obraz:

docker pull debezium/connect

Minimalny zestaw zmiennych środowiskowych wymaganych do uruchomienia konektora wygląda następująco:

  • BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092,kafka-3:9092 — początkowa lista serwerów klastra Kafka dla uzyskania pełnej listy członków klastra;
  • OFFSET_STORAGE_TOPIC=connector-offsets — topic do przechowywania pozycji, na której obecnie znajduje się konektor;
  • CONNECT_STATUS_STORAGE_TOPIC=connector-status — topic do przechowywania statusu konektora i jego zadań;
  • CONFIG_STORAGE_TOPIC=connector-config — topic do przechowywania danych konfiguracyjnych konektora i jego zadań;
  • GROUP_ID=1 — identyfikator grupy workerów, na których może być wykonywane zadanie konektora; wymagany przy użyciu trybu rozproszonego (distributed) trybu.

Uruchamiamy kontener z tymi zmiennymi:

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

Uwaga na temat Avro

Domyślnie Debezium zapisuje dane w formacie JSON, co jest akceptowalne dla środowisk testowych i niewielkich ilości danych, ale może stać się problemem w wysokoobciążonych bazach. Alternatywą dla konwertera JSON jest serializacja wiadomości przy użyciu Avro do formatu binarnego, co pozwala na zmniejszenie obciążenia podsystemu I/O w Apache Kafka.

Aby korzystać z Avro, należy wdrożyć oddzielny schema-registry (do przechowywania schematów). Zmienne dla konwertera będą wyglądały następująco:

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

Szczegóły dotyczące korzystania z Avro i konfigurowania rejestru wykraczają poza zakres artykułu — w dalszej części dla klarowności będziemy używać JSON.

2. Konfiguracja samego konektora

Teraz możemy przejść bezpośrednio do konfiguracji samego konektora, który będzie odczytywał dane ze źródła.

Rozważymy na przykładzie konektorów dla dwóch DBMS: PostgreSQL i MongoDB — z którymi mam doświadczenie i które mają różnice (nawet jeśli drobne, w niektórych przypadkach — istotne!).

Konfiguracja jest opisana w notacji JSON i ładowana do Kafka Connect za pomocą żądania POST.

2.1. PostgreSQL

Przykład konfiguracji konektora dla 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"
  }
}

Zasada działania konektora po takiej konfiguracji jest dość prosta:

  • Przy pierwszym uruchomieniu łączy się z bazą, określoną w konfiguracji, i przechodzi w tryb initial snapshot, wysyłając do Kafka początkowy zbiór danych, uzyskanych za pomocą warunku SELECT * FROM table_name.
  • Po zakończeniu inicjalizacji konektor przechodzi w tryb odczytu zmian z plików WAL PostgreSQL.

O dostępnych opcjach:

  • name — nazwa konektora, dla którego używana jest konfiguracja opisana poniżej; ta nazwa jest później wykorzystywana do zarządzania konektorem (tj. przeglądanie statusu / ponowne uruchamianie / aktualizowanie konfiguracji) przez REST API Kafka Connect;
  • connector.class — klasa konektora DBMS, która będzie używana przez skonfigurowany konektor;
  • plugin.name — nazwa wtyczki do logicznego dekodowania danych z plików WAL. Do wyboru dostępne są wal2json, decoderbuffs i pgoutput. Pierwsze dwie wymagają zainstalowania odpowiednich rozszerzeń w DBMS, a pgoutput dla PostgreSQL w wersji 10 i wyższej nie wymagają dodatkowych manipulacji;
  • database.* — opcje do połączenia z DB, gdzie database.server.name — nazwa instancji PostgreSQL, używana do tworzenia nazwy tematu w klastrze Kafka;
  • table.include.list — lista tabel, które chcemy monitorować pod kątem zmian; określa się w formacie schema.table_name; nie można używać razem z table.exclude.list;
  • heartbeat.interval.ms — interwał (w milisekundach), z jakim konektor wysyła wiadomości heartbeat do specjalnego tematu;
  • heartbeat.action.query — zapytanie, które będzie wykonywane przy wysyłaniu każdej wiadomości heartbeat (opcja dostępna od wersji 1.1);
  • slot.name — nazwa slotu replikacji, który będzie używany przez konektor;
  • publication.name — nazwa publikacji w PostgreSQL, której używa konektor. W przypadku, gdy nie istnieje, Debezium spróbuje ją utworzyć. Jeśli użytkownik, z którym następuje połączenie, nie ma wystarczających uprawnień do tego działania — konektor zakończy działanie z błędem;
  • transforms określa, jak należy zmieniać nazwę celu tematu:
    • transforms.AddPrefix.type określa, że będziemy używać wyrażeń regularnych;
    • transforms.AddPrefix.regex — maska, według której nadpisywana jest nazwa celu tematu;
    • transforms.AddPrefix.replacement — bezpośrednio to, na co nadpisujemy.

Więcej o heartbeat i transforms

Domyślnie konektor wysyła dane do Kafka przy każdym zatwierdzeniu transakcji, a jego LSN (Log Sequence Number) zapisuje w temacie serwisowym offset. Ale co się stanie, jeśli konektor jest skonfigurowany do odczytu nie całej bazy, a tylko części jej tabel (w których aktualizacja danych nie zachodzi często)?

  • Konektor będzie odczytywał pliki WAL i nie wykryje w nich zatwierdzenia transakcji w tych tabelach, które monitoruje.
  • Dlatego nie będzie aktualizował swojej bieżącej pozycji ani w temacie, ani w slocie replikacji.
  • To z kolei doprowadzi do „zatrzymania” plików WAL na dysku i prawdopodobnego wyczerpania całej przestrzeni dyskowej.

I tutaj na pomoc przychodzą opcje heartbeat.interval.ms i heartbeat.action.query. Użycie tych opcji w parze pozwala za każdym razem przy wysyłaniu wiadomości heartbeat wykonywać zapytanie o zmianę danych w osobnej tabeli. Dzięki temu stale aktualizowany jest LSN, na którym obecnie znajduje się konektor (w slocie replikacji). To pozwala systemowi bazodanowemu usunąć pliki WAL, które nie są już potrzebne. Więcej informacji na temat działania opcji można znaleźć w dokumentacji.

Inna opcja, która zasługuje na większą uwagę, to transforms. Chociaż bardziej chodzi o wygodę i estetykę…

Domyślnie Debezium tworzy tematy, kierując się następującą polityką nazewnictwa: serverName.schemaName.tableName. To nie zawsze może być wygodne. Opcje transforms pozwalają określać listę tabel, dla których zdarzenia należy kierować do tematu o konkretnej nazwie, za pomocą wyrażeń regularnych.

W naszej konfiguracji dzięki transforms dzieje się co następuje: wszystkie zdarzenia CDC z monitorowanej bazy danych trafią do tematu o nazwie data.cdc.dbname. W przeciwnym wypadku (bez tych ustawień) Debezium domyślnie tworzyłby temat dla każdej tabeli w postaci: pg-dev.public.

.

Ograniczenia konektora

Na zakończenie opisu konfiguracji konektora dla PostgreSQL warto wspomnieć o następujących cechach/ograniczeniach jego działania:

  1. Funkcjonalność konektora dla PostgreSQL opiera się na koncepcji logicznego dekodowania. Dlatego nie śledzi on zapytań o zmianę struktury bazy danych (DDL) — w związku z tym, w tematach tych danych nie będzie.
  2. Ponieważ używane są sloty replikacji, połączenie konektora jest możliwe tylko z wiodącym instancją DBMS.
  3. Jeśli użytkownik, pod którym konektor łączy się z bazą danych, ma przyznane tylko prawa do odczytu, to przed pierwszym uruchomieniem konieczne będzie ręczne utworzenie slotu replikacji i publikacji w bazie danych.

Zastosowanie konfiguracji

Zatem załadujemy naszą konfigurację do konektora:

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

Sprawdzamy, czy załadunek przebiegł pomyślnie i konektor został uruchomiony:

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

Świetnie: jest skonfigurowany i gotowy do pracy. Teraz udajemy się do odbiorcy i łączymy się z Kafka, a następnie dodajemy i zmieniamy wpis w tabeli:

$ 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

W naszym temacie odbije się to w następujący sposób:

Bardzo długi JSON z naszymi zmianami

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

W obu przypadkach zapisy składają się z klucza (PK) rekordu, który został zmieniony, oraz samej istoty zmian: jaki był rekord przed i jaki jest po zmianie.

  • W przypadku INSERT: wartość przed (before) równa się null, a po — ciąg, który został wstawiony.
  • W przypadku UPDATE: w payload.before wyświetlane jest poprzednie stanu wiersza, a w payload.after — nowym z istotą zmian.

2.2 MongoDB

Ten konektor wykorzystuje standardowy mechanizm replikacji MongoDB, odczytując informacje z oplog'a węzła primary bazy danych.

Podobnie jak opisany wcześniej konektor dla PgSQL, tutaj również przy pierwszym uruchomieniu wykonywana jest podstawowa migawka danych, po czym konektor przechodzi w tryb odczytu z oplog'a.

Przykład konfiguracji:

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

Jak można zauważyć, brakuje tutaj nowych opcji w porównaniu do poprzedniego przykładu, jednak zmniejszyła się tylko liczba opcji dotyczących połączenia z bazą danych i ich prefiksów.

Ustawienia transforms tym razem robią to, co następuje: przekształcają nazwę docelowego topiku z wzoru .. do data.cdc.mongo_.

Odporność na awarie

Kwestię odporności i wysokiej dostępności rozwiązuje się obecnie jak nigdy dotąd — zwłaszcza gdy mówimy o danych i transakcjach, a monitorowanie zmian danych również wpisuje się w te zagadnienia. Rozważmy, co może pójść nie tak i co będzie się działo z Debezium w każdym z przypadków.

Są trzy warianty awarii:

  1. Awaria Kafka Connect. Jeśli Connect jest skonfigurowany do pracy w trybie rozproszonym, należy kilku pracownikom nadać ten sam group.id. Wtedy w przypadku awarii jednego z nich konektor zostanie uruchomiony na innym pracowniku i będzie kontynuował odczyt z ostatniej zatwierdzonej pozycji w topiku w Kafka.
  2. Utrata łączności z klastrem Kafka. Konektor po prostu zatrzyma odczyt na pozycji, której nie udało się przesłać do Kafka, i okresowo będzie próbował ponownego jej przesłania, aż do momentu, gdy próba zakończy się sukcesem.
  3. Niedostępność źródła danych. Connector będzie próbował ponownego połączenia z źródłem zgodnie z konfiguracją. Domyślnie to 16 prób z wykorzystaniem exponential backoff. Po 16. nieudanej próbie zadanie zostanie oznaczone jako failed i będzie wymagać ręcznego ponownego uruchomienia przez interfejs REST Kafka Connect.
    • W przypadku PostgreSQL Dane nie znikną, ponieważ użycie slotów replikacji uniemożliwi usunięcie plików WAL, które nie zostały odczytane przez connector. W tym przypadku istnieje jednak ciemna strona: jeśli na dłuższy czas zostanie przerwana łączność sieciowa między connector a DBMS, istnieje ryzyko, że miejsce na dysku się skończy, co może prowadzić do całkowitego awarii DBMS.
    • W przypadku MySQL Pliki binlogów mogą zostać obrócone przez sam DBMS wcześniej, niż przywróci się łączność. To spowoduje, że connector przejdzie w stan failed, a do przywrócenia normalnego działania będzie wymagać ponownego uruchomienia w trybie initial snapshot, aby kontynuować odczyt z binlogów.
    • O MongoDB. Dokumentacja stwierdza: zachowanie connectora w przypadku, gdy pliki dziennika/ oplog zostały usunięte i connector nie może kontynuować odczytu z pozycji, w której się zatrzymał, jest takie samo dla wszystkich DBMS. Polega to na tym, że connector przejdzie w stan failed i będzie wymagać ponownego uruchomienia w trybie initial snapshot.

      Jednak zdarzają się wyjątki. Jeśli connector był wyłączony przez dłuższy czas (lub nie miał dostępu do instancji MongoDB), a oplog w tym czasie przeszedł rotację, to po przywróceniu połączenia connector spokojnie zacznie odczytywać dane od pierwszej dostępnej pozycji, przez co część danych w Kafka nie zostanie utracona.

Podsumowanie

Debezium — moje pierwsze doświadczenie z systemami CDC i ogólnie bardzo pozytywne. Projekt przyciąga wsparciem dla głównych DBMS, prostotą konfiguracji, wsparciem dla klasteryzacji i aktywną społecznością. Zainteresowanym praktyką polecam zapoznać się z przewodnikami dla Kafka Connect i Debezium.

W porównaniu do JDBC Connector dla Kafka Connect, główną zaletą Debezium jest to, że zmiany są odczytywane z dzienników baz danych, co pozwala na pozyskiwanie danych z minimalnym opóźnieniem. JDBC Connector (dostarczany z Kafka Connect) wykonuje zapytania do monitorowanej tabeli w stałych odstępach czasu i (z tego samego powodu) nie generuje wiadomości przy usuwaniu danych (jak można zapytać o dane, których nie ma?).

Aby rozwiązać podobne zadania, warto zwrócić uwagę na następujące rozwiązania (oprócz Debezium):

P.S.

Przeczytaj także na naszym blogu:

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster