W przededniu uruchomienia nowego kursu przygotowaliśmy tłumaczenie interesującego materiału.

Przegląd
Omówimy dość popularny wzorzec, w ramach którego aplikacje wykorzystują kilka magazynów danych, z których każdy jest używany do swoich celów, na przykład do przechowywania kanonicznej formy danych (MySQL itd.), zapewnienia rozszerzonych możliwości wyszukiwania (ElasticSearch itd.), pamięci podręcznej (Memcached itd.) i innych. Zwykle przy użyciu kilku magazynów danych jedno z nich działa jako główny magazyn, a pozostałe jako podrzędne. Jedynym problemem jest to, jak synchronizować te magazyny danych.
Przyjrzeliśmy się kilku różnym wzorcom, które próbowały rozwiązać problem synchronizacji kilku magazynów, takim jak podwójne zapisy, transakcje rozproszone itd. Jednak te podejścia mają istotne ograniczenia w zakresie użyteczności w rzeczywistości, niezawodności i konserwacji. Oprócz synchronizacji danych, niektóre aplikacje muszą także wzbogacać dane, wywołując zewnętrzne usługi.
W celu rozwiązania tych problemów stworzono Delta. Delta jest ostatecznie spójną, zarządzaną wydarzeniami platformą do synchronizacji i wzbogacania danych.
Istniejące rozwiązania
Podwójne zapisy
Aby zsynchronizować dwa magazyny danych, można użyć podwójnego zapisu, który najpierw dokonuje zapisu w jednym magazynie, a zaraz po tym zapisuje w drugim. Pierwszy zapis można powtórzyć, a drugi przerwać, jeśli pierwszy zakończy się niepowodzeniem po wyczerpaniu prób. Jednak dwa magazyny danych mogą przestać się synchronizować, jeśli zapis w drugim magazynie zakończy się niepowodzeniem. Problem ten zwykle rozwiązuje się przez stworzenie procedury naprawczej, która może okresowo przenosić dane z pierwszego magazynu do drugiego lub czynić to tylko w przypadku wykrycia różnic w danych.
Problemy:
Wykonanie procedury przywracania to specyficzne zadanie, które nie może być ponownie używane. Ponadto dane między magazynami pozostają niesynchroniczne, dopóki nie zakończy się procedura przywracania. Problemy nastręcza to, jeśli używa się więcej niż dwóch magazynów danych. W końcu procedura przywracania może zwiększyć obciążenie pierwotnego źródła danych.
Tabela logów zmian
Gdy w zestawie tabel zachodzą zmiany (na przykład wstawianie, aktualizacja i usuwanie wpisów), wpisy zmian są dodawane do tabeli logów jako część tej samej transakcji. Inny wątek lub proces nieustannie żąda zdarzeń z tabeli logów i zapisuje je w jednym lub kilku magazynach danych, usuwając zdarzenia z tabeli logów po potwierdzeniu zapisania przez wszystkie magazyny.
Problemy:
Ten wzorzec powinien być realizowany jako biblioteka, najlepiej bez zmiany kodu aplikacji, która go używa. W środowisku poliglotycznym realizacja takiej biblioteki powinna istnieć w każdym potrzebnym języku, ale zapewnienie spójności funkcji i zachowań między językami jest bardzo trudne.
Inny problem polega na uzyskaniu zmian schematu w systemach, które nie obsługują transakcyjnych zmian schematu [1][2], takich jak MySQL. Dlatego wzór wykonywania zmiany (na przykład zmiany schematu) oraz transakcyjnego zapisania go w tabeli logów zmian nie zawsze będzie działać.
Transakcje Rozproszone
Transakcje rozproszone mogą być używane do rozdzielenia transakcji między różne heterogeniczne magazyny danych, tak aby operacja została zatwierdzona we wszystkich używanych magazynach lub nie została zatwierdzona w żadnym z nich.
Problemy:
Rozproszone transakcje to bardzo duży problem dla heterogenicznych magazynów danych. Zasadniczo mogą one polegać tylko na najmniejszym wspólnym mianowniku uczestniczących systemów. Na przykład transakcje XA blokują wykonanie, gdy w trakcie aplikacji wystąpi błąd na etapie przygotowania. Ponadto XA nie zapewnia wykrywania zakleszczeń i nie obsługuje optymistycznych schematów zarządzania równoległością. Ponadto niektóre systemy, takie jak ElasticSearch, nie obsługują XA ani żadnego innego heterogenicznego modelu transakcji. W związku z tym zapewnienie atomowości zapisu w różnych technologiach przechowywania danych pozostaje dla aplikacji bardzo trudnym zadaniem [3].
Delta
Delta została stworzona w celu rozwiązania ograniczeń istniejących rozwiązań synchronizacji danych, a także umożliwia wzbogacanie danych w czasie rzeczywistym. Naszym celem było abstrahowanie wszystkich tych skomplikowanych kwestii od programistów aplikacji, aby mogli w pełni skupić się na wdrażaniu funkcjonalności biznesowej. Poniżej opiszemy „Movie Search”, rzeczywisty przypadek użycia Delta w Netflixie.
W Netflixie powszechnie stosuje się architekturę mikroserwisową, a każdy mikroserwis zwykle obsługuje jeden typ danych. Podstawowe informacje o filmie zostały przeniesione do mikroserwisu o nazwie Movie Service, a związane z nimi dane, takie jak informacje o producentach, aktorach, dostawcach itd., są zarządzane przez kilka innych mikroserwisów (konkretniej Deal Service, Talent Service i Vendor Service).
Użytkownicy biznesowi w Netflix Studios często potrzebują wyszukiwać filmy według różnych kryteriów, dlatego ważne jest dla nich, aby mieć możliwość przeszukiwania wszystkich danych związanych z filmami.
Przed wprowadzeniem Delta zespół wyszukiwania filmów musiał pozyskiwać dane z kilku mikroserwisów, zanim zindeksowałby dane o filmach. Dodatkowo zespół musiał opracować system, który okresowo aktualizowałby indeks wyszukiwania, żądając zmian od innych mikroserwisów, nawet jeśli nie było żadnych zmian. System szybko stał się skomplikowany i trudny do utrzymania.

Rysunek 1. System pollingu przed wprowadzeniem Delta
Po rozpoczęciu korzystania z Delta, system został uproszczony do systemu zarządzanego zdarzeniami, jak pokazano na poniższym rysunku. Zdarzenia CDC (Change-Data-Capture) są wysyłane do tematów Keystone Kafka za pomocą Delta-Connector. Aplikacja Delta, zbudowana z użyciem Delta Stream Processing Framework (opartego na Flink), otrzymuje zdarzenia CDC z tematu, wzbogaca je, wywołując inne mikroserwisy, a ostatecznie przesyła wzbogacone dane do indeksu wyszukiwania w Elasticsearch. Cały proces odbywa się niemal w czasie rzeczywistym, co oznacza, że jak tylko zmiany zostaną zarejestrowane w hurtowni danych, indeksy wyszukiwania są aktualizowane.

Rysunek 2. Pipeline danych przy użyciu Delta
W następnych rozdziałach opiszemy działanie Delta-Connector, który łączy się z hurtownią danych i publikuje zdarzenia CDC na poziomie transportu, który stanowi infrastrukturę do przesyłania danych w czasie rzeczywistym, kierując zdarzenia CDC do tematów Kafka. Na końcu porozmawiamy o strukturze przetwarzania strumieni Delta, którą mogą wykorzystać programiści aplikacji do logiki przetwarzania i wzbogacania danych.
CDC (Change-Data-Capture)
Opracowaliśmy usługę CDC o nazwie Delta-Connector, która może rejestrować zatwierdzone zmiany z hurtowni danych w czasie rzeczywistym i zapisywać je w strumieniu. Zmiany w czasie rzeczywistym pochodzą z dzienników transakcji i zrzutów hurtowni. Zrzuty są używane, ponieważ dzienniki transakcji zazwyczaj nie przechowują pełnej historii zmian. Zmiany są zazwyczaj serializowane jako zdarzenia Delta, dzięki czemu odbiorca nie musi się martwić o źródło zmiany.
Delta-Connector obsługuje kilka dodatkowych funkcji, takich jak:
- Możliwość zapisywania w niestandardowych wyjściach z pominięciem Kafka.
- Możliwość uruchamiania ręcznych zrzutów w dowolnym momencie dla wszystkich tabel, określonej tabeli lub dla konkretnych kluczy głównych.
- Zrzuty można pobierać w kawałkach, więc nie ma potrzeby rozpoczynania od początku w przypadku awarii.
- Nie ma potrzeby blokowania tabel, co jest bardzo ważne, aby ruch zapisu do bazy danych nigdy nie był blokowany przez naszą usługę.
- Wysoka dostępność dzięki zapasowym instancjom w strefach dostępności AWS.
Obecnie obsługujemy MySQL i Postgres, w tym przy wdrażaniu w AWS RDS i Aurora. Obsługujemy również Cassandra (multi-master). Więcej szczegółów dotyczących Delta-Connector można znaleźć w tym .
Kafka i poziom transportu
Poziom transportu zdarzeń Delta oparty jest na usłudze wymiany wiadomości platformy .
Historycznie, publikacja wiadomości w Netflix była optymalizowana pod kątem zwiększenia dostępności, a nie trwałości (patrz ). Kompromisem było potencjalne niespójność danych brokera w różnych skrajnych scenariuszach. Na przykład, unclean leader election jest odpowiedzialny za to, że odbiorca potencjalnie duplikuje lub traci zdarzenia.
Z Delta chcieliśmy uzyskać bardziej solidne gwarancje trwałości, aby zapewnić dostarczanie zdarzeń CDC do pochodnych magazynów danych. W tym celu zaproponowaliśmy specjalnie zaprojektowany klaster Kafka jako obiekt pierwszej klasy. Możesz zobaczyć niektóre ustawienia brokera w tabeli poniżej:

W klastrach Keystone Kafka, unclean leader election zazwyczaj włączony w celu zapewnienia dostępności wydawcy. Może to prowadzić do utraty wiadomości, jeśli niesynchronizowana replika zostanie wybrana jako lider. Dla nowego, wysoko niezawodnego klastra Kafka parametr unclean leader election wyłączony, aby zapobiec utracie wiadomości.
Zwiększyliśmy również replication factor z 2 do 3 i minimum insync replicas z 1 do 2. Wydawcy, piszący do tego klastra, wymagają acks od wszystkich innych, gwarantując, że 2 z 3 replik będą miały najnowsze wiadomości wysłane przez wydawcę.
Kiedy instancja brokera kończy działanie, nowa instancja zastępuje starą. Jednak nowemu brokerowi będzie musiał nadrobić niezsynchonizowane repliki, co może zająć kilka godzin. Aby skrócić czas przywracania działania tego scenariusza, zaczęliśmy używać magazynu blokowego (Amazon Elastic Block Store) zamiast lokalnych dysków brokerów. Gdy nowa instancja zastępuje zakończoną instancję brokera, podłącza wolumen EBS, który był używany przez zakończoną instancję, i zaczyna doganiać nowe wiadomości. Ten proces skraca czas likwidacji opóźnienia z kilku godzin do kilku minut, ponieważ nowej instancji nie trzeba już replikować z pustego stanu. Ogólnie rzecz biorąc, oddzielne cykle życia magazynu i brokera znacznie redukują wpływ efektu zmiany brokera.
Aby jeszcze bardziej zwiększyć gwarancję dostarczenia danych, użyliśmy w celu wykrycia jakiejkolwiek utraty wiadomości w ekstremalnych warunkach (na przykład, rozsynchornizowania zegarów w liderze podziału).
Stream Processing Framework
Poziom przetwarzania w Delta oparty jest na platformie Netflix SPaaS, która zapewnia integrację Apache Flink z ekosystemem Netflix. Platforma udostępnia interfejs użytkownika, który zarządza wdrażaniem zadań Flink i orkiestracją klastrów Flink nad naszą platformą zarządzania kontenerami Titus. Interfejs ten również zarządza konfiguracjami zadań i pozwala użytkownikom wprowadzać zmiany w konfiguracji dynamicznie, bez konieczności ponownej kompilacji zadań Flink.
Delta dostarcza framework do przetwarzania strumieniowego (stream processing framework) danych oparty na Flink i SPaaS, który wykorzystuje oparty na adnotacjach DSL (Domain Specific Language), aby zatuszować techniczne szczegóły. Na przykład, aby zdefiniować krok, według którego będą wzbogacane zdarzenia, wywołując zewnętrzne usługi, użytkownicy muszą napisać następujący DSL, a framework stworzy na jego podstawie model, który będzie działał w Flink.

Rysunek 3. Przykład wzbogacania na DSL w Delta
Framework do przetwarzania nie tylko skraca krzywą uczenia się, ale również zapewnia wspólne funkcje przetwarzania strumieniowego, takie jak deduplikacja, schematyzacja, a także elastyczność i odporność na awarie dla rozwiązywania typowych problemów w pracy.
Delta Stream Processing Framework składa się z dwóch kluczowych modułów: modułu DSL & API oraz modułu Runtime. Moduł DSL & API zapewnia DSL i UDF (User-Defined-Function) API, aby użytkownicy mogli napisać własną logikę przetwarzania (na przykład filtrowanie lub transformacje). Moduł Runtime zapewnia implementację parsera DSL, który buduje wewnętrzną reprezentację kroków przetwarzania w modelach DAG. Komponent Execution interpretuje modele DAG, aby zainicjować rzeczywiste operatory Flink i ostatecznie uruchomić aplikację Flink. Architektura frameworka została zilustrowana na poniższym rysunku.

Rysunek 4. Architektura Delta Stream Processing Framework
Takie podejście ma kilka zalet:
- Użytkownicy mogą skupić się na swojej logice biznesowej bez konieczności zagłębiania się w specyfikę Flink lub strukturę SPaaS.
- Optymalizacja może być realizowana w sposób przejrzysty dla użytkowników, a błędy mogą być naprawiane bez konieczności wprowadzania jakichkolwiek zmian w kodzie użytkownika (UDF).
- Praca aplikacji Delta została uproszczona dla użytkowników, ponieważ platforma zapewnia elastyczność i odporność na błędy od podstaw oraz zbiera wiele szczegółowych metryk, które można wykorzystać do powiadomień.
Użycie w produkcji
Delta działa w produkcji od ponad roku i odgrywa kluczową rolę w wielu aplikacjach Netflix Studio. Pomogła zespołom wdrożyć takie zastosowania jak indeksowanie wyszukiwania, przechowywanie danych i przepływy pracy zarządzane zdarzeniami. Poniżej przedstawiono przegląd wysokopoziomowej architektury platformy Delta.

Rysunek 5. Wysokopoziomowa architektura Delta.
Podziękowania
Chcielibyśmy podziękować następującym osobom, które uczestniczyły w tworzeniu i rozwoju Delta w Netflix: Allen Wang, Charles Zhao, Jaebin Yoon, Josh Snyder, Kasturi Chatterjee, Mark Cho, Olof Johansson, Piyush Goyal, Prashanth Ramdas, Raghuram Onti Srinivasan, Sandeep Gupta, Steven Wu, Tharanga Gamaethige, Yun Wang i Zhenzhong Xu.
Źródła
- Martin Kleppmann, Alastair R. Beresford, Boerge Svingen: Przetwarzanie zdarzeń w czasie rzeczywistym. Commun. ACM 62(5): 43–49 (2019). DOI:
: „Data Build Tool dla magazynu Amazon Redshift”.
Źródło: habr.com
