Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Cześć, Habr!

Przypominamy, że po książce o Kafka opublikowaliśmy nie mniej interesującą pracę o bibliotece Kafka Streams API.

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Na razie społeczność dopiero odkrywa granice możliwości tego potężnego narzędzia. Niedawno ukazał się artykuł, z tłumaczeniem którego chcemy Państwa zapoznać. Autor opisuje na swoim doświadczeniu, jak stworzyć z Kafka Streams rozproszone magazyn danych. Miłej lektury!

Biblioteka Apache Kafka Streams jest używana na całym świecie w przedsiębiorstwach do rozproszonego przetwarzania strumieniowego na bazie Apache Kafka. Jednym z niedocenianych aspektów tego frameworka jest to, że pozwala na utrzymywanie lokalnego stanu, generowanego na podstawie przetwarzania strumieniowego.

W tym artykule opowiem, jak w naszej firmie udało się efektywnie wykorzystać tę możliwość przy opracowywaniu produktu dotyczącego bezpieczeństwa aplikacji w chmurze. Przy pomocy Kafka Streams stworzyliśmy mikroserwisy z dzielonym stanem, z których każdy jest dla nas odpornością na awarie oraz wysoko dostępnym źródłem wiarygodnych informacji o stanie obiektów w systemie. To krok naprzód zarówno pod względem niezawodności, jak i wygody wsparcia.

Jeśli interesuje Państwa alternatywne podejście, które umożliwia wykorzystanie jednego centralnego bazy danych do wsparcia formalnego stanu Państwa obiektów – przeczytajcie, warto…

Dlaczego uznaliśmy, że nadszedł czas na zmianę naszych podejść do pracy z dzielonym stanem

Musieliśmy utrzymywać stan różnych obiektów, opierając się na raportach agentów (np. czy strona była atakowana?). Przed przejściem na Kafka Streams często polegaliśmy na zarządzaniu stanem za pomocą jednej centralnej bazy danych (+ API usługowe). Takie podejście ma swoje wady: w sytuacjach intensywnych danych utrzymanie spójności i synchronizacji staje się prawdziwym wyzwaniem. Baza danych może stać się wąskim gardłem lub znajdować się w stanie wyścigu i cierpieć z powodu nieprzewidywalności.

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 1: typowy scenariusz z podziałem stanu, które występowały przed przejściem na
Kafka i Kafka Streams: agenci przekazują swoje reprezentacje przez API, zaktualizowany stan jest obliczany za pomocą centralnej bazy danych.

Poznajcie Kafka Streams – tworzenie mikroserwisów z dzielonym stanem stało się proste.

Około rok temu postanowiliśmy dokładnie przeanalizować nasze scenariusze pracy z dzielonym stanem, aby rozwiązać pewne problemy. Od razu postanowiliśmy wypróbować Kafka Streams – wiadomo, jak bardzo jest skalowalna, wysoko dostępna i odporna na błędy oraz jak bogate ma funkcjonalności strumieniowe (przekształcenia, w tym z zachowaniem stanu). Dokładnie to, czego potrzebowaliśmy, nie wspominając o tym, jak dojrzałym i niezawodnym systemem wymiany wiadomości jest Kafka.

Każdy z naszych stworzonych mikroserwisów z zachowaniem stanu oparty był na instancji Kafka Streams z dość prostą topologią. Składała się ona z 1) źródła 2) procesora z trwałym magazynem kluczy i wartości 3) odcinka:

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 2: domyślna topologia naszych instancji strumieniowych dla mikroserwisów z zachowaniem stanu. Zauważ, że znajduje się tutaj także magazyn, w którym przetrzymywane są metadane dotyczące harmonogramowania.

W tym nowym podejściu agenci tworzą wiadomości wprowadzane do początkowego tematu, a konsumenci – na przykład, usługa powiadomień e-mail – odbierają obliczony dzielony stan przez odcinek (temat wyjściowy).

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 3: nowy przykład strumienia zadań dla scenariusza z dzielonymi mikroserwisami: 1) agent generuje wiadomość, która trafia do początkowego tematu Kafka; 2) mikroserwis z dzielonym stanem (używający Kafka Streams) przetwarza ją i zapisuje obliczony stan w końcowym temacie Kafka; po czym 3) konsumenci odbierają nowy stan.

Hej, a to wbudowane magazyn kluczy i wartości jest naprawdę bardzo przydatne!

Jak wspomniano wcześniej, nasza topologia z dzielonym stanem zawiera magazyn kluczy i wartości. Znaleźliśmy kilka sposobów jego wykorzystania, z których dwa opisane są poniżej.

Opcja #1: wykorzystanie magazynu kluczy i wartości w obliczeniach.

Nasze pierwsze magazyn kluczy i wartości zawierał dane pomocnicze potrzebne do naszych obliczeń. Na przykład, w niektórych przypadkach współdzielony stan był określany na zasadzie „większości głosów”. W magazynie można było przechowywać wszystkie ostatnie raporty agentów dotyczące stanu pewnego obiektu. Następnie, otrzymując nowy raport od danego agenta, mogliśmy go zapisać, wydobyć z magazynu raporty wszystkich innych agentów dotyczące tego samego obiektu i powtórzyć obliczenia.
Na poniższej ilustracji 4 pokazano, jak otwieraliśmy dostęp do magazynu kluczy i wartości metodzie przetwarzającej procesora, dzięki czemu można było przetworzyć nową wiadomość.

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 4: otwieramy dostęp do magazynu kluczy i wartości dla przetwarzającej metody procesora (po tym każdy scenariusz operujący na współdzielonym stanie musi zaimplementować metodę doProcess)

Opcja #2: tworzenie API CRUD na bazie Kafka Streams

Po ustabilizowaniu naszego podstawowego strumienia zadań zaczęliśmy próbować napisać RESTful CRUD API dla naszych mikroserwisów z współdzielonym stanem. Chcieliśmy, aby można było pobierać stan niektórych lub wszystkich obiektów, a także ustawiać lub usuwać stan obiektu (co jest przydatne w przypadku utrzymania backendu).

Aby zapewnić wsparcie dla wszystkich API Get State, za każdym razem, gdy musieliśmy ponownie obliczać stan podczas przetwarzania, przez długi czas przechowywaliśmy go w wbudowanym magazynie kluczy i wartości. W takim przypadku wystarczy zaimplementować takie API za pomocą jednego wystąpienia Kafka Streams, jak pokazano w poniższym kodzie:

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 5: użycie wbudowanego magazynu kluczy i wartości do uzyskania wcześniej obliczonego stanu obiektu

Aktualizacja stanu obiektu przez API również jest łatwa do zrealizowania. W zasadzie, aby to zrobić, trzeba tylko stworzyć producenta Kafka, a za jego pomocą dokonać zapisu zawierającego nowy stan. Dzięki temu zapewniono, że wszystkie wiadomości generowane przez API będą przetwarzane dokładnie w taki sam sposób, jak te pochodzące od innych producentów (np. agentów).

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 6: stan obiektu można ustawić za pomocą producenta Kafka

Małe komplikacje: Kafka ma wiele partycji

Następnie chcieliśmy rozłożyć obciążenie związane z przetwarzaniem i poprawić dostępność, dostarczając dla każdego scenariusza klaster mikroserwisów ze współdzielonym stanem. Konfiguracja okazała się dla nas dziecinnie prosta: po skonfigurowaniu wszystkich instancji tak, aby pracowały z tym samym ID aplikacji (i tymi samymi serwerami początkowymi), praktycznie wszystko pozostałe odbywało się automatycznie. Ustaliliśmy również, że każdy temat źródłowy będzie składał się z kilku partycji, aby każdej instancji przypisać podzbiór takich partycji.

Warto również wspomnieć, że tutaj w porządku rzeczy jest tworzenie kopii zapasowej magazynu stanów, aby na przykład w przypadku przywracania po awarii przenieść tę kopię na inną instancję. Dla każdego magazynu stanów w Kafka Streams tworzony jest replicable topic z dziennikiem zmian (w którym śledzone są lokalne aktualizacje). W ten sposób Kafka ciągle zabezpiecza magazyn stanów. Dlatego w przypadku awarii danej instancji, magazyn stanów Kafka Streams może być szybko przywrócony na innej instancji, na którą przeniosą się odpowiednie partycje. Nasze testy pokazały, że trwa to zaledwie kilka sekund, nawet jeśli w magazynie znajdują się miliony wpisów.

Przechodząc od jednego mikroserwisu ze współdzielonym stanem do klastra mikroserwisów, wdrożenie Get State API nie jest już tak trywialne. W nowej sytuacji w magazynie stanów każdego mikroserwisu znajduje się tylko część ogólnego obrazu (te obiekty, których klucze były mapowane na konkretną partycję). Musieliśmy określić, na której instancji znajdował się stan poszukiwanego obiektu, co robiliśmy na podstawie metadanych strumieni, jak pokazano poniżej:

Nie tylko przetwarzanie: Jak zrobiliśmy z Kafka Streams rozproszoną bazę danych i co z tego wynikło

Ilustracja 7: przy pomocy metadanych strumieni określamy, z której instancji żądać stanu poszukiwanego obiektu; podobne podejście stosowane było z GET ALL API.

Główne wnioski

Magazyny stanów w Kafka Streams de facto mogą służyć jako rozproszona baza danych,

  • ciągle replikowana w Kafka.
  • Na takiej systemie łatwo zbudować API CRUD.
  • Przetwarzanie wielu partycji staje się nieco bardziej skomplikowane.
  • Możliwe jest również dodanie jednego lub więcej magazynów stanów do topologii strumieniowej w celu przechowywania danych pomocniczych. Taka opcja może być używana do:
  • Długoterminowego przechowywania danych potrzebnych do obliczeń w czasie przetwarzania strumieniowego
  • Długoterminowego przechowywania danych, które mogą być przydatne podczas następnej inicjalizacji instancji strumienia
  • wielu innych…

Dzięki tym i innym zaletom Kafka Streams doskonale nadaje się do wspierania globalnego stanu w tak rozproszonej systemie jak nasz. Kafka Streams okazała się bardzo niezawodna w produkcji (od momentu jej wdrożenia praktycznie nie utraciliśmy wiadomości), a jesteśmy pewni, że jej możliwości nie są na tym ograniczone!

Ź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