Delta: Piattaforma di sincronizzazione dei dati e arricchimento

In vista del lancio di un nuovo ciclo del corso «Data Engineer» abbiamo preparato la traduzione di un materiale interessante.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento

Panoramica

Parleremo di un pattern piuttosto popolare che consente alle applicazioni di utilizzare diversi archivi di dati, ciascuno dei quali è impiegato per scopi specifici, come la memorizzazione della forma canonica dei dati (MySQL, ecc.), l'offerta di capacità di ricerca avanzate (ElasticSearch, ecc.), la memorizzazione nella cache (Memcached, ecc.) e altro. Di solito, nell'uso di più archivi di dati, uno di essi funge da archivio principale, mentre gli altri come archivi secondari. L'unico problema è come sincronizzare questi archivi di dati.

Abbiamo esaminato diversi pattern che cercano di affrontare il problema della sincronizzazione di più archivi, come la doppia scrittura, le transazioni distribuite, ecc. Tuttavia, questi approcci presentano notevoli limiti in termini di utilizzo nella vita reale, affidabilità e manutenzione tecnica. Oltre alla sincronizzazione dei dati, alcune applicazioni devono anche arricchire i dati richiedendo servizi esterni.

Per risolvere questi problemi è stata sviluppata Delta. Delta rappresenta alla fine una piattaforma coerente, gestita dagli eventi, per la sincronizzazione e l'arricchimento dei dati.

Soluzioni esistenti

Doppia scrittura

Per sincronizzare due archivi di dati, si può utilizzare la doppia scrittura, che esegue una scrittura in un archivio e subito dopo una scrittura nell'altro. La prima scrittura può essere ripetuta, mentre la seconda può essere interrotta se la prima fallisce dopo un numero limitato di tentativi. Tuttavia, i due archivi di dati possono smettere di sincronizzarsi se la scrittura nel secondo archivio fallisce. Questo problema è solitamente risolto creando una procedura di recupero che può periodicamente trasferire nuovamente i dati dal primo archivio al secondo o farlo solo se si riscontrano differenze nei dati.

Problemi:

L'esecuzione della procedura di recupero è un'operazione specifica che non può essere riutilizzata. Inoltre, i dati tra gli archivi rimangono desincronizzati fino a quando non viene completata la procedura di recupero. La situazione si complica se si utilizzano più di due archivi di dati. Infine, la procedura di recupero può aumentare il carico sull'origine dati iniziale.

Tabella dei log delle modifiche

Quando si verificano modifiche in un insieme di tabelle (ad esempio, inserimenti, aggiornamenti e cancellazioni di record), le registrazioni delle modifiche vengono aggiunte a una tabella di log come parte della stessa transazione. Un altro thread o processo interroga costantemente gli eventi dalla tabella di log e li scrive in uno o più archivi di dati, eliminando eventualmente gli eventi dalla tabella di log dopo che il loro inserimento è stato confermato da tutti gli archivi.

Problemi:

Questo pattern dovrebbe essere implementato come una libreria e, idealmente, senza modificare il codice dell'applicazione che la utilizza. In un ambiente poliglotta, l'implementazione di tale libreria dovrebbe esistere in qualsiasi linguaggio necessario, ma garantire la coerenza del funzionamento delle funzioni e dei comportamenti tra i linguaggi risulta molto complicato.

Un'altra problematica riguarda l'acquisizione delle modifiche nello schema, in quei sistemi che non supportano le modifiche di schema transazionali [1][2], come ad esempio MySQL. Pertanto, il pattern di esecuzione del cambiamento (ad esempio, le modifiche allo schema) e la registrazione transazionale nelle tabelle di log delle modifiche non sempre funzionerà.

Transazioni distribuite

Le transazioni distribuite possono essere utilizzate per dividere una transazione tra diversi archivi di dati eterogenei, in modo che l'operazione venga registrata in tutti gli archivi utilizzati oppure non venga registrata in nessuno di essi.

Problemi:

Le transazioni distribuite rappresentano un problema rilevante per archivi di dati eterogenei. Per loro natura possono contare solo sul minimo comune denominatore dei sistemi partecipanti. Ad esempio, le transazioni XA bloccano l'esecuzione se si verifica un errore durante la fase di preparazione. Inoltre, XA non fornisce la rilevazione dei deadlock né supporta schemi di gestione della concorrenza ottimistica. Inoltre, alcuni sistemi come ElasticSearch non supportano XA o qualsiasi altro modello di transazione eterogeneo. Pertanto, garantire l'atomicità della registrazione in diverse tecnologie di archiviazione rimane una sfida complessa per le applicazioni [3].

Delta

Delta è stata progettata per superare i limiti delle soluzioni esistenti per la sincronizzazione dei dati e consente anche di arricchire i dati in tempo reale. Il nostro obiettivo era astraRe tutti questi aspetti complessi agli sviluppatori, in modo che potessero concentrarsi completamente sull'implementazione delle funzionalità aziendali. Di seguito descriveremo "Movie Search", il caso d'uso reale di Delta da Netflix.

In Netflix si utilizza ampiamente un'architettura a microservizi, con ogni microservizio che generalmente gestisce un solo tipo di dato. I dettagli principali sui film sono gestiti da un microservizio chiamato Movie Service, mentre i dati correlati, come le informazioni sui produttori, sugli attori e sui fornitori, sono gestiti da diversi altri microservizi (in particolare Deal Service, Talent Service e Vendor Service).
Gli utenti aziendali di Netflix Studios spesso necessitano di cercare film in base a diversi criteri, motivo per cui è fondamentale avere la possibilità di cercare in tutti i dati relativi ai film.

Prima dell'introduzione di Delta, il team di ricerca film doveva raccogliere dati da diversi microservizi prima di indicizzare i dati sui film. Inoltre, il team doveva sviluppare un sistema che aggiornasse periodicamente l'indice di ricerca richiedendo modifiche agli altri microservizi, anche quando non c'erano modifiche. Questo sistema è rapidamente diventato complesso e difficile da mantenere.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Figura 1. Sistema di polling prima di Delta
Dopo l'implementazione di Delta, il sistema è stato semplificato in un sistema basato su eventi, come mostrato nella figura seguente. Gli eventi CDC (Change-Data-Capture) vengono inviati ai topic di Keystone Kafka tramite Delta-Connector. L'applicazione Delta, costruita utilizzando il Delta Stream Processing Framework (basato su Flink), riceve eventi CDC dal topic, li arricchisce chiamando altri microservizi e, infine, trasmette i dati arricchiti all'indice di ricerca in Elasticsearch. L'intero processo avviene quasi in tempo reale, ovvero, non appena le modifiche vengono registrate nel data store, gli indici di ricerca vengono aggiornati.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Figura 2. Pipeline dei dati con Delta
Nei prossimi capitoli descriveremo il funzionamento di Delta-Connector, che si collega al data store e pubblica eventi CDC a livello di trasporto, che rappresenta l'infrastruttura di trasmissione dati in tempo reale, dirigendo gli eventi CDC verso i topic Kafka. Alla fine parleremo della struttura di elaborazione dei flussi Delta, che gli sviluppatori possono utilizzare per la logica di elaborazione e arricchimento dei dati.

CDC (Change-Data-Capture)

Abbiamo sviluppato un servizio CDC chiamato Delta-Connector, che può rilevare modifiche confermate dal data store in tempo reale e scriverle in un flusso. Le modifiche in tempo reale provengono dal registro delle transazioni e dai dump dello storage. I dump sono usati poiché i registri delle transazioni di solito non conservano l'intera cronologia delle modifiche. Le modifiche vengono solitamente serializzate come eventi Delta, così il destinatario non deve preoccuparsi di dove origina la modifica.

Delta-Connector supporta diverse funzionalità aggiuntive, come:

  • La possibilità di scrivere in output personalizzati al di fuori di Kafka.
  • La possibilità di attivare dump manuali in qualsiasi momento per tutte le tabelle, per una specifica tabella o per determinate chiavi primarie.
  • I dump possono essere recuperati a blocchi, quindi non è necessario ricominciare dall'inizio in caso di errore.
  • Non è necessario bloccare le tabelle, il che è molto importante affinché il traffico di scrittura nel database non venga mai bloccato dal nostro servizio.
  • Alta disponibilità grazie a istanze di backup nelle AWS Availability Zones.

Attualmente supportiamo MySQL e Postgres, incluso il deployment su AWS RDS e Aurora. Supportiamo anche Cassandra (multi-master). Maggiori dettagli su Delta-Connector possono essere trovati in questo blog.

Kafka e livello di trasporto

Il livello di trasporto degli eventi Delta è costruito su un servizio di messaggistica della piattaforma. Keystone.

Storicamente, la pubblicazione di messaggi in Netflix è stata ottimizzata per aumentare la disponibilità, piuttosto che la durabilità (vedi l'articolo precedente). Il compromesso è stato il potenziale disallineamento dei dati broker in vari scenari limite. Ad esempio, l'elezione sporca del leader può causare la duplicazione o la perdita di eventi da parte del destinatario.

Con Delta volevamo ottenere garanzie più solide di durata, per garantire la consegna di eventi CDC negli archivi derivati. A tal fine, abbiamo proposto un cluster Kafka progettato su misura come oggetto di prima classe. Puoi vedere alcune impostazioni del broker nella tabella qui sotto:

Delta: Piattaforma di sincronizzazione dei dati e arricchimento

Nei cluster Keystone Kafka, l'elezione sporca del leader è solitamente attivato per garantire la disponibilità del publisher. Questo può portare a una perdita di messaggi se una replica non sincronizzata viene scelta come leader. Per un nuovo cluster Kafka ad alta affidabilità, l'impostazione l'elezione sporca del leader è disattivata per prevenire la perdita di messaggi.

Inoltre, abbiamo aumentato il fattore di replica da 2 a 3 e il numero minimo di repliche in sync da 1 a 2. I publisher che scrivono in questo cluster richiedono acks da tutte le altre repliche, garantendo che 2 su 3 repliche avranno i messaggi più recenti inviati dal publisher.

Quando un'istanza del broker si arresta, una nuova istanza sostituisce quella vecchia. Tuttavia, al nuovo broker sarà necessario allinearsi con le repliche non sincronizzate, il che può richiedere alcune ore. Per ridurre il tempo di recupero in questo scenario, abbiamo iniziato a utilizzare lo storage a blocchi (Amazon Elastic Block Store) invece dei dischi locali dei broker. Quando una nuova istanza sostituisce un'istanza di broker terminata, si collega al volume EBS che era dell'istanza terminata e inizia a sincronizzarsi con i nuovi messaggi. Questo processo riduce il tempo di recupero da diverse ore a pochi minuti, poiché la nuova istanza non deve più replicare da uno stato vuoto. In generale, i cicli di vita separati dello storage e del broker riducono significativamente l'impatto dell'effetto di cambio broker.

Per aumentare ulteriormente la garanzia di consegna dei dati, abbiamo utilizzato un sistema di tracciamento dei messaggi per rilevare qualsiasi perdita di messaggi in condizioni estreme (ad esempio, la dissincronizzazione degli orari nel leader della partizione).

Stream Processing Framework

Il livello di elaborazione in Delta è costruito sulla piattaforma Netflix SPaaS, che fornisce integrazione con Apache Flink nell'ecosistema Netflix. La piattaforma offre un'interfaccia utente che gestisce il deployment dei job Flink e l'orchestrazione dei cluster Flink sulla nostra piattaforma di gestione dei container Titus. L'interfaccia gestisce anche le configurazioni dei job e consente agli utenti di apportare modifiche dinamiche senza dover ricompilare i job Flink.

Delta fornisce un framework di elaborazione dei dati in streaming basato su Flink e SPaaS, che utilizza un DSL (Domain Specific Language) basato su annotazioni, per astrarre i dettagli tecnici. Ad esempio, per definire il passaggio in cui arricchire gli eventi, chiamando servizi esterni, gli utenti devono scrivere il seguente DSL, e il framework creerà un modello che verrà eseguito su Flink. Figura 3. Esempio di arricchimento su DSL in Delta

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Il framework di elaborazione non solo riduce la curva di apprendimento, ma offre anche funzioni comuni di elaborazione dei flussi, come la deduplicazione, la schematizzazione, oltre a flessibilità e resilienza per affrontare comuni problemi operativi.

Delta Stream Processing Framework è composto da due moduli chiave, il modulo DSL & API e il modulo Runtime. Il modulo DSL & API fornisce un API DSL e UDF (User-Defined-Function) affinché gli utenti possano scrivere la propria logica di elaborazione (ad esempio, filtri o trasformazioni). Il modulo Runtime fornisce un'implementazione del parser DSL che costruisce una rappresentazione interna dei passi di elaborazione in modelli DAG. Il componente Execution interpreta i modelli DAG per inizializzare i reali operatori Flink e infine avviare l'applicazione Flink. L'architettura del framework è illustrata nella figura seguente.

Figura 4. Architettura del Delta Stream Processing Framework

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Questo approccio ha diversi vantaggi:

Gli utenti possono concentrarsi sulla propria logica aziendale senza dover entrare nei dettagli specifici di Flink o nella struttura di SPaaS.

  • L'ottimizzazione può essere eseguita in modo trasparente per gli utenti, e gli errori possono essere risolti senza la necessità di modifiche nel codice utente (UDF).
  • L'operatività delle applicazioni Delta è semplificata per gli utenti, poiché la piattaforma offre flessibilità e resilienza già pronte all'uso e raccoglie una serie di metriche dettagliate che possono essere utilizzate per allerta.
  • Utilizzo in produzione

Utilizzo in produzione

Delta è in produzione da oltre un anno e gioca un ruolo chiave in molte applicazioni di Netflix Studio. Ha supportato i team nell'implementazione di casi d'uso come indicizzazione della ricerca, archiviazione dei dati e flussi di lavoro gestiti da eventi. Di seguito è riportata una panoramica dell'architettura ad alto livello della piattaforma Delta.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Figura 5. Architettura ad alto livello di Delta.

Ringraziamenti

Vorremmo ringraziare le seguenti persone che hanno contribuito alla creazione e allo sviluppo di Delta in 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 e Zhenzhong Xu.

Fonti

  1. dev.mysql.com/doc/refman/5.7/en/implicit-commit.html
  2. dev.mysql.com/doc/refman/5.7/en/cannot-roll-back.html
  3. Martin Kleppmann, Alastair R. Beresford, Boerge Svingen: Elaborazione di eventi online. Commun. ACM 62(5): 43–49 (2019). DOI: doi.org/10.1145/3312527

Registrati per il webinar gratuito: «Data Build Tool per il data warehousing su Amazon Redshift».

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS 🔥 Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS | ProHoster