Im Vorfeld des Starts des neuen Kurses Wir haben eine interessante Materialübersetzung vorbereitet.

Überblick
Wir werden über ein recht populäres Muster sprechen, das Anwendungen verwenden, um mehrere Datenspeicher zu nutzen, wobei jeder Speicher für seine eigenen Zwecke verwendet wird, z.B. zur Speicherung der kanonischen Form von Daten (MySQL usw.), zur Bereitstellung erweiterter Suchmöglichkeiten (ElasticSearch usw.), zur Zwischenspeicherung (Memcached usw.) und anderen. In der Regel arbeitet eines der Datenspeicher als primäres Speicher und die anderen als derivative Speicher. Das einzige Problem besteht darin, wie man diese Datenspeicher synchronisiert.
Wir haben eine Reihe verschiedener Muster betrachtet, die versucht haben, das Problem der Synchronisierung mehrerer Speicher zu lösen, wie z.B. Double-Write, verteilte Transaktionen usw. Diese Ansätze haben jedoch wesentliche Einschränkungen hinsichtlich der Anwendung in der Praxis, der Zuverlässigkeit und der Wartung. Neben der Synchronisierung von Daten müssen einige Anwendungen auch Daten anreichern, indem sie externe Dienste aufrufen.
Um diese Probleme zu lösen, wurde Delta entwickelt. Delta stellt letztendlich eine konsistente, ereignisgesteuerte Plattform zur Synchronisierung und Anreicherung von Daten dar.
Bestehende Lösungen
Double-Write
Um zwei Datenspeicher zu synchronisieren, kann der Double-Write-Mechanismus verwendet werden, der einen Schreibvorgang in einem Speicher ausführt und dann sofort danach im anderen Speicher. Der erste Schreibvorgang kann wiederholt werden, während der zweite abgebrochen werden kann, wenn der erste nach Erschöpfung der Versuche fehlschlägt. Allerdings können sich zwei Datenspeicher unsynchronisieren, wenn der Schreibvorgang im zweiten Speicher fehlschlägt. Dieses Problem wird in der Regel durch die Einrichtung eines Wiederherstellungsverfahrens gelöst, das periodisch Daten vom ersten Speicher in den zweiten übertragen kann oder dies nur dann tut, wenn Unterschiede in den Daten festgestellt werden.
Probleme:
Die Durchführung des Wiederherstellungsverfahrens ist eine spezifische Arbeit, die nicht wiederverwendet werden kann. Darüber hinaus bleiben die Daten zwischen den Speicherorten unsynchronisiert, bis das Wiederherstellungsverfahren abgeschlossen ist. Die Situation wird komplizierter, wenn mehr als zwei Datenspeicher verwendet werden. Schließlich kann das Wiederherstellungsverfahren die Last auf die ursprüngliche Datenquelle erhöhen.
Änderungsprotokolltabelle
Wenn es in der Tabellensammlung Änderungen gibt (z. B. Einfügen, Aktualisieren und Löschen von Datensätzen), werden die Änderungen in die Änderungsprotokolltabelle als Teil derselben Transaktion eingefügt. Ein anderer Thread oder Prozess fragt ständig Ereignisse aus der Änderungsprotokolltabelle ab und schreibt sie in einen oder mehrere Datenspeicher, wobei die Ereignisse nach Bestätigung des Eintrags durch alle Speicher entfernt werden, falls notwendig.
Probleme:
Dieses Muster sollte als Bibliothek implementiert werden und idealerweise ohne Änderungen am verwendenden Anwendungs-Code. In einer Polyglotsprache sollte eine solche Bibliothek in jeder erforderlichen Sprache existieren, jedoch ist es sehr schwierig, die Konsistenz der Funktionalitäten und des Verhaltens zwischen den Sprachen zu gewährleisten.
Ein weiteres Problem liegt in der Erfassung von Schemaänderungen in Systemen, die keine transaktionalen Schemaänderungen unterstützen [1][2], wie zum Beispiel MySQL. Daher wird das Muster zur Durchführung von Änderungen (z. B. Schemaänderungen) und der transaktionalen Aufzeichnung in die Änderungsprotokolltabelle nicht immer funktionieren.
Verteilte Transaktionen
Verteilte Transaktionen können verwendet werden, um eine Transaktion zwischen mehreren heterogenen Datenspeichern zu teilen, sodass die Operation entweder in allen verwendeten Speichern festgelegt oder in keinem von ihnen festgelegt wird.
Probleme:
Verteilte Transaktionen sind ein gravierendes Problem für heterogene Datenspeicher. Ihrer Natur nach können sie sich nur auf den kleinsten gemeinsamen Nenner der beteiligten Systeme stützen. Zum Beispiel blockieren XA-Transaktionen die Ausführung, wenn während des Anwendungsprozesses ein Fehler in der Vorbereitungsphase auftritt. Darüber hinaus bietet XA keine Erkennung von Deadlocks und unterstützt keine optimistischen Parallelitätssteuerungsschemata. Außerdem unterstützen einige Systeme wie ElasticSearch kein XA oder ein anderes heterogenes Transaktionsmodell. Daher bleibt die Gewährleistung der Atomarität von Aufzeichnungen in verschiedenen Datenspeichertechnologien eine sehr komplexe Aufgabe für Anwendungen [3].
Delta
Delta wurde entwickelt, um die Einschränkungen bestehender Lösungen zur Datensynchronisierung zu beseitigen, und ermöglicht es zudem, Daten in Echtzeit anzureichern. Unser Ziel war es, all diese komplexen Aspekte von den Anwendungsentwicklern zu abstrahieren, damit sie sich voll und ganz auf die Umsetzung der Geschäftslogik konzentrieren können. Im Folgenden werden wir "Movie Search" beschreiben, den tatsächlichen Anwendungsfall von Delta bei Netflix.
Bei Netflix wird weitgehend eine Microservice-Architektur verwendet, und jeder Microservice bedient normalerweise einen bestimmten Datentyp. Die grundlegenden Informationen zu Filmen werden in einen Microservice namens Movie Service ausgelagert, während verwandte Daten wie Informationen über Produzenten, Schauspieler, Anbieter usw. von mehreren anderen Microservices verwaltet werden (namentlich Deal Service, Talent Service und Vendor Service).
Die Business-Nutzer in Netflix Studios müssen oft Filme nach verschiedenen Kriterien durchsuchen, weshalb es für sie von großer Bedeutung ist, die Möglichkeit zu haben, alle filmbezogenen Daten zu durchsuchen.
Vor der Einführung von Delta musste das Film-Suchteam Daten aus mehreren Microservices abrufen, bevor sie die Filmdaten indizierten. Darüber hinaus musste das Team ein System entwickeln, das den Suchindex regelmäßig aktualisierte, indem es bei anderen Microservices nach Änderungen fragte, selbst wenn es überhaupt keine Änderungen gab. Dieses System geriet schnell in Komplexität und wurde schwer zu warten.

Abbildung 1. Das Polling-System vor Delta
Nach dem Start der Verwendung von Delta wurde das System auf ein ereignisgesteuertes System vereinfacht, wie im folgenden Bild gezeigt. CDC-Ereignisse (Change-Data-Capture) werden über den Delta-Connector in die Themen von Keystone Kafka gesendet. Die Delta-Anwendung, die mit dem Delta Stream Processing Framework (basierend auf Flink) erstellt wurde, erhält CDC-Ereignisse aus dem Thema, bereichert sie durch Aufrufe an andere Mikrodienste und überträgt schließlich die angereicherten Daten in den Suchindex in Elasticsearch. Der gesamte Prozess erfolgt nahezu in Echtzeit, das heißt, sobald Änderungen im Datenspeicher erfasst werden, werden die Suchindizes aktualisiert.

Abbildung 2. Datenpipeline bei Verwendung von Delta
In den folgenden Abschnitten werden wir die Funktionsweise des Delta-Connectors beschreiben, der sich mit dem Datenspeicher verbindet und CDC-Ereignisse auf Transportschicht veröffentlicht, die die Infrastruktur für die Übertragung von Daten in Echtzeit darstellen, die CDC-Ereignisse in Kafka-Themen leitet. Am Ende werden wir über die Struktur der Delta-Streaming-Verarbeitung sprechen, die Entwickler für die Logik der Datenverarbeitung und -anreicherung verwenden können.
CDC (Change-Data-Capture)
Wir haben einen CDC-Dienst namens Delta-Connector entwickelt, der in der Lage ist, in Echtzeit committete Änderungen aus dem Datenspeicher zu erfassen und in einen Stream zu schreiben. Echtzeitänderungen stammen aus dem Transaktionsprotokoll und den Dumps des Datenspeichers. Dumps werden verwendet, da Transaktionsprotokolle normalerweise nicht die gesamte Historie der Änderungen speichern. Änderungen werden üblicherweise als Delta-Ereignisse serialisiert, sodass sich der Empfänger keine Sorgen machen muss, woher die Änderung kommt.
Der Delta-Connector unterstützt mehrere zusätzliche Funktionen, wie zum Beispiel:
- Die Möglichkeit, direkt in benutzerdefinierte Ausgaben neben Kafka zu schreiben.
- Die Möglichkeit, manuelle Dumps jederzeit für alle Tabellen, eine bestimmte Tabelle oder für bestimmte Primärschlüssel zu aktivieren.
- Dumps können in Chunks abgerufen werden, sodass es nicht notwendig ist, im Falle eines Fehlers von vorne zu beginnen.
- Es ist nicht erforderlich, Sperren auf Tabellen zu setzen, was sehr wichtig ist, damit der Schreibverkehr in die Datenbank niemals durch unseren Dienst blockiert wird.
- Hohe Verfügbarkeit durch redundante Instanzen in AWS Availability Zones.
Derzeit unterstützen wir MySQL und Postgres, einschließlich des Einsatzes in AWS RDS und Aurora. Außerdem unterstützen wir Cassandra (multi-master). Weitere Informationen zum Delta-Connector finden Sie in diesem .
Kafka und Transportschicht
Die Transportschicht der Delta-Ereignisse basiert auf dem Messaging-Service der Plattform .
Historisch gesehen wurde die Veröffentlichung von Nachrichten in Netflix auf die Erhöhung der Verfügbarkeit und nicht auf die Langlebigkeit optimiert (siehe ). Ein Kompromiss war eine potenzielle Inkonsistenz der Broker-Daten in verschiedenen Grenzfällen. Zum Beispiel, unclean leader election verantwortet dafür, dass der Empfänger potenziell Ereignisse dupliziert oder verliert.
Mit Delta wollten wir robustere Garantien für die Langlebigkeit erhalten, um die Lieferung von CDC-Ereignissen an abgeleitete Speicherorte sicherzustellen. Zu diesem Zweck haben wir einen speziell gestalteten Kafka-Cluster als First-Class-Objekt vorgeschlagen. Sie können einige Broker-Einstellungen in der Tabelle unten ansehen:

In den Keystone Kafka-Clustern, unclean leader election ist normalerweise aktiviert, um die Verfügbarkeit des Publishers zu gewährleisten. Dies kann zu Nachrichtenverlust führen, wenn eine nicht synchronisierte Replikation als Leader gewählt wird. Für den neuen hochverfügbaren Kafka-Cluster ist die Einstellung unclean leader election deaktiviert, um Nachrichtenverlust zu verhindern.
Außerdem haben wir den Replikationsfaktor von 2 auf 3 und die minimalen insync Replicas von 1 auf 2 erhöht. Publisher, die in diesen Cluster schreiben, benötigen Acks von allen anderen, um sicherzustellen, dass 2 von 3 Replikaten die aktuellsten Nachrichten des Publishers erhalten.
Wenn eine Brokerinstanz beendet wird, ersetzt eine neue Instanz die alte. Allerdings muss der neue Broker die nicht synchronisierten Replikate aufholen, was mehrere Stunden in Anspruch nehmen kann. Um die Wiederherstellungszeit dieses Szenarios zu verkürzen, haben wir begonnen, Blockdatenspeicher (Amazon Elastic Block Store) anstelle von lokalen Festplatten der Broker zu verwenden. Wenn die neue Instanz die beendete Brokerinstanz ersetzt, schließt sie das EBS-Laufwerk an, das der beendeten Instanz zugeordnet war, und beginnt, neue Nachrichten aufzuholen. Dieser Prozess reduziert die Zeit zur Beseitigung des Rückstands von mehreren Stunden auf wenige Minuten, da die neue Instanz nicht mehr aus einem leeren Zustand replizieren muss. Insgesamt minimieren die separaten Lebenszyklen von Speicher und Broker die Auswirkungen des Brokerwechsels erheblich.
Um die Datenliefergarantie weiter zu erhöhen, haben wir implementiert, um jeglichen Nachrichtenverlust unter extremen Bedingungen (z. B. Zeitdifferenzen zwischen den Knoten eines Partitionführers) zu erkennen.
Stream Processing Framework
Die Verarbeitungsstufe in Delta basiert auf der Netflix SPaaS-Plattform, die Apache Flink mit dem Netflix-Ökosystem integriert. Die Plattform bietet eine Benutzeroberfläche, die die Bereitstellung von Flink-Jobs und die Orchestrierung von Flink-Clustern über unsere Container-Management-Plattform Titus verwaltet. Die Benutzeroberfläche verwaltet auch die Jobkonfigurationen und ermöglicht es den Benutzern, Änderungen an der Konfiguration dynamisch vorzunehmen, ohne die Flink-Jobs neu kompilieren zu müssen.
Delta bietet ein Framework für die Verarbeitung von Streamingdaten (stream processing framework) auf der Grundlage von Flink und SPaaS, die annotationsbasiert DSL (Domain Specific Language) verwendet, um technische Details zu abstrahieren. Um beispielsweise zu definieren, wie die Ereignisse angereichert werden sollen, indem externe Dienste aufgerufen werden, müssen die Benutzer das folgende DSL schreiben, und das Framework erstellt basierend darauf ein Modell, das durch Flink ausgeführt wird.

Abbildung 3. Beispiel für Anreicherung in DSL in Delta
Das Framework zur Verarbeitung reduziert nicht nur die Lernkurve, sondern bietet auch allgemeine Funktionen für die Stromverarbeitung, wie Duplizierungsschutz, Schematierung sowie Flexibilität und Fehlertoleranz zur Lösung allgemeiner Probleme im Betrieb.
Das Delta Stream Processing Framework besteht aus zwei Hauptmodulen: einem DSL- & API-Modul und einem Runtime-Modul. Das DSL- & API-Modul bietet eine DSL und eine UDF (User-Defined Function) API, damit Benutzer ihre eigenen Verarbeitungslogiken (zum Beispiel Filterung oder Transformationen) erstellen können. Das Runtime-Modul stellt eine Implementierung des DSL-Parsers bereit, der eine interne Darstellung der Verarbeitungsschritte in DAG-Modellen erstellt. Die Execution-Komponente interpretiert die DAG-Modelle, um die tatsächlichen Flink-Operatoren zu initialisieren und letztendlich die Flink-Anwendung zu starten. Die Architektur des Frameworks wird im folgenden Bild illustriert.

Abb. 4. Architektur des Delta Stream Processing Frameworks
Dieser Ansatz bietet mehrere Vorteile:
- Benutzer können sich auf ihre Geschäftslogik konzentrieren, ohne sich mit den Feinheiten von Flink oder der SPaaS-Struktur auseinandersetzen zu müssen.
- Optimierungen können für die Benutzer transparent durchgeführt werden, und Fehler können behoben werden, ohne dass Änderungen am Benutzercode (UDF) erforderlich sind.
- Die Arbeit mit Delta-Anwendungen wird für die Benutzer erleichtert, da die Plattform von Haus aus Flexibilität und Fehlertoleranz bietet und zahlreiche detaillierte Metriken erfasst, die für Benachrichtigungen verwendet werden können.
Einsatz in der Produktion
Delta wird bereits seit über einem Jahr in der Produktion eingesetzt und spielt eine Schlüsselrolle in vielen Anwendungen des Netflix Studios. Es hat den Teams geholfen, Anwendungsfälle wie Suchindexierung, Datenspeicherung und ereignisgesteuerte Workflows zu realisieren. Im Folgenden wird eine Übersicht über die hochgradige Architektur der Delta-Plattform gegeben.

Abb. 5. Hochgradige Architektur von Delta.
Dank
Wir möchten den folgenden Personen danken, die an der Entwicklung und dem Wachstum von Delta bei Netflix beteiligt waren: 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 und Zhenzhong Xu.
Quellen
- Martin Kleppmann, Alastair R. Beresford, Boerge Svingen: Online Event Processing. Commun. ACM 62(5): 43–49 (2019). DOI:
: „Data Build Tool für das Amazon Redshift Data Warehouse“.
Quelle: habr.com
