Delta: Piattaforma di sincronizzazione dei dati e arricchimento

In attesa del lancio di un nuovo flusso relativo al corso «Data Engineer» abbiamo preparato la traduzione di un materiale interessante.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento

Panoramica

Parleremo di un modello abbastanza popolare, attraverso il quale le applicazioni utilizzano più archivi di dati, dove ciascun archivio è utilizzato per i propri scopi, ad esempio, per memorizzare la forma canonica dei dati (MySQL, ecc.), per garantire capacità di ricerca avanzate (ElasticSearch, ecc.), per la memorizzazione temporanea (Memcached, ecc.) e altro. Di solito, quando si utilizza più di un archivio di dati, uno di essi funge da archivio principale e gli altri da archivi secondari. L'unico problema è come sincronizzare questi archivi di dati.

Abbiamo esaminato una serie di modelli diversi che cercavano di risolvere il problema della sincronizzazione di più archivi, come la scrittura doppia, transazioni distribuite, ecc. Tuttavia, questi approcci presentano limitazioni significative in termini di utilizzo nella vita reale, affidabilità e manutenzione. Oltre alla sincronizzazione dei dati, alcune applicazioni necessitano anche di arricchire i dati, chiamando servizi esterni.

Per risolvere questi problemi è stata sviluppata Delta. Delta rappresenta in ultima analisi una piattaforma coerente, gestita da eventi per la sincronizzazione e l'arricchimento dei dati.

Soluzioni esistenti

Scrittura doppia

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

Problemi:

L'esecuzione della procedura di ripristino è un lavoro specifico che non può essere riutilizzato. Inoltre, i dati tra i repository rimangono disallineati fino al termine della procedura di ripristino. La situazione si complica se si utilizzano più di due repository di dati. Infine, la procedura di ripristino può aggiungere un carico alla fonte di dati originale.

Tabella dei log delle modifiche

Quando si verificano modifiche nel set di tabelle (ad esempio, inserimento, aggiornamento e cancellazione di una registrazione), le registrazioni delle modifiche vengono aggiunte alla tabella dei log come parte della stessa transazione. Un altro thread o processo interroga costantemente gli eventi dalla tabella dei log e li scrive in uno o più repository di dati, eliminando gli eventi dalla tabella dei log dopo che tutti i repository hanno confermato la scrittura.

Problemi:

Questo modello 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 del comportamento tra i linguaggi è molto complicato.

Un'altra problematica riguarda l'ottenimento delle modifiche allo schema, in quei sistemi che non supportano le modifiche transazionali allo schema [1][2], come ad esempio MySQL. Pertanto, il modello di esecuzione della modifica (ad esempio, modifiche allo schema) e la registrazione transazionale nella tabella dei log delle modifiche non funzionano sempre.

Transazioni Distribuite

Le transazioni distribuite possono essere utilizzate per suddividere una transazione tra più repository di dati eterogenei, in modo tale che l'operazione venga registrata in tutti i repository utilizzati oppure non venga registrata in nessuno di essi.

Problemi:

Le transazioni distribuite sono un problema molto grande per i dati eterogenei. Per loro natura, possono fare affidamento 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 nel processo dell'applicazione. Inoltre, XA non fornisce rilevazione dei deadlock e non supporta schemi di gestione della concorrenza ottimistica. Oltre a ciò, alcuni sistemi come ElasticSearch non supportano XA o qualsiasi altro modello di transazione eterogeneo. Pertanto, garantire l'atomicità della scrittura in varie tecnologie di archiviazione dei dati rimane una sfida molto complessa per le applicazioni [3].

Delta

Delta è stata sviluppata per superare i limiti delle soluzioni esistenti per la sincronizzazione dei dati e permette anche di arricchire i dati in tempo reale. Il nostro obiettivo era astrarre tutte queste complessità dagli sviluppatori di applicazioni, affinché potessero concentrarsi completamente sulla realizzazione delle funzionalità di business. Successivamente descriveremo «Movie Search», il caso d'uso effettivo di Delta da parte di Netflix.

In Netflix si utilizza ampiamente un'architettura a microservizi e ogni microservizio di solito gestisce un solo tipo di dati. Le informazioni di base sui film sono gestite in un microservizio chiamato Movie Service, mentre i dati correlati, come le informazioni su produttori, attori, fornitori e così via, sono gestiti da diversi altri microservizi (ovvero Deal Service, Talent Service e Vendor Service).
Gli utenti business di Netflix Studios hanno spesso bisogno di cercare film in base a diversi criteri, motivo per cui è essenziale per loro avere la possibilità di cercare in tutti i dati relativi ai film.

Prima dell'arrivo di Delta, il team di ricerca dei film doveva raccogliere i dati da diversi microservizi prima di indicizzare i dati dei 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 state 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'inizio dell'utilizzo di Delta, il sistema è stato semplificato in un sistema basato su eventi, come mostrato nell'immagine seguente. Gli eventi CDC (Change-Data-Capture) vengono inviati nei topic di Keystone Kafka tramite Delta-Connector. L'applicazione Delta, costruita utilizzando il Delta Stream Processing Framework (basato su Flink), riceve gli 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 warehouse, gli indici di ricerca vengono aggiornati.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Figura 2. Pipeline dei dati durante l'uso di Delta
Nei prossimi capitoli descriveremo il funzionamento del Delta-Connector, che si connette al data warehouse e pubblica gli eventi CDC a livello di trasporto, che rappresenta l'infrastruttura di trasmissione dei dati in tempo reale, indirizzando gli eventi CDC nei topic Kafka. E infine parleremo della struttura di elaborazione dei flussi Delta, che gli sviluppatori delle applicazioni 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ò registrare le modifiche commesse dal data warehouse in tempo reale e scriverle in un flusso. Le modifiche in tempo reale vengono tratte dal registro delle transazioni e dai dump del data warehouse. I dump vengono utilizzati perché i registri delle transazioni di solito non conservano l'intera cronologia delle modifiche. Le modifiche vengono generalmente serializzate come eventi Delta, in modo che il destinatario non debba preoccuparsi da dove provenga la modifica.

Delta-Connector supporta diverse funzioni aggiuntive, come:

  • La possibilità di scrivere in output personalizzati bypassando Kafka.
  • La possibilità di attivare dump manuali in qualsiasi momento per tutte le tabelle, una tabella specifica o per determinate chiavi primarie.
  • I dump possono essere recuperati a blocchi, quindi non è necessario ricominciare da capo 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.
  • Elevata disponibilità grazie a istanze di backup nelle AWS Availability Zones.

Attualmente supportiamo MySQL e Postgres, compresi nel 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 sul 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 un potenziale disallineamento dei dati del broker in vari scenari di confine. Ad esempio, elezione del leader non pulita responsabile del fatto che il destinatario potenzialmente duplica o perde eventi.

Con Delta volevamo ottenere garanzie più solide sulla durabilità per garantire la consegna degli eventi CDC nei magazzini derivati. A tal fine, abbiamo proposto un cluster Kafka progettato specificamente come oggetto di prima classe. Puoi visualizzare alcune impostazioni del broker nella tabella sottostante:

Delta: Piattaforma di sincronizzazione dei dati e arricchimento

Nei cluster Keystone Kafka, elezione del leader non pulita di solito è abilitato per garantire la disponibilità del publisher. Ciò può portare alla perdita di messaggi se una replica non sincronizzata viene scelta come leader. Per un nuovo cluster Kafka ad alta affidabilità, l'impostazione elezione del leader non pulita è disabilitata per prevenire la perdita di messaggi.

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

Quando un'istanza del broker termina, una nuova istanza sostituisce quella vecchia. Tuttavia, il nuovo broker dovrà recuperare le repliche non sincronizzate, il che può richiedere diverse ore. Per ridurre il tempo di ripristino di questo scenario, abbiamo iniziato a utilizzare lo storage a blocchi (Amazon Elastic Block Store) al posto dei dischi locali dei broker. Quando una nuova istanza sostituisce un'istanza del broker che ha terminato, essa collega il volume EBS che era dell'istanza terminata e inizia a recuperare i nuovi messaggi. Questo processo riduce il tempo necessario per eliminare il ritardo da diverse ore a pochi minuti, poiché al nuovo broker non è più necessario replicare dallo stato vuoto. In generale, i cicli di vita separati dello storage e del broker riducono significativamente l'impatto del cambio del 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, desincronizzazione degli orari nel leader della partizione).

Stream Processing Framework

Il livello di elaborazione in Delta si basa sulla piattaforma Netflix SPaaS, che integra Apache Flink nell'ecosistema Netflix. La piattaforma fornisce un'interfaccia utente che gestisce il deployment dei job Flink e l'orchestrazione dei cluster Flink sopra la nostra piattaforma di gestione dei container Titus. L'interfaccia gestisce anche le configurazioni dei job e consente agli utenti di apportare modifiche alle configurazioni in modo dinamico senza dover ricompilare i job Flink.

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

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Il framework di elaborazione non solo accorcia la curva di apprendimento, ma garantisce anche funzionalità comuni per l'elaborazione dei flussi, come la deduplicazione, la schematizzazione, oltre a flessibilità e resilienza per affrontare problemi comuni nelle operazioni.

Il framework per l'elaborazione non solo riduce la curva di apprendimento, ma fornisce anche funzionalità comuni di elaborazione dei flussi, come la deduplicazione, la schematizzazione, oltre a flessibilità e tolleranza ai guasti per affrontare i problemi comuni nel lavoro.

Il 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 nei modelli DAG. Il componente Execution interpreta i modelli DAG per inizializzare i veri operatori Flink e infine avviare l'applicazione Flink. L'architettura del framework è illustrata nella figura seguente.

Delta: Piattaforma di sincronizzazione dei dati e arricchimento
Figura 4. Architettura del Delta Stream Processing Framework

Questo approccio presenta diversi vantaggi:

  • Gli utenti possono concentrarsi sulla loro logica di business senza dover approfondire le specifiche di Flink o la struttura di SPaaS.
  • L'ottimizzazione può avvenire in modo trasparente per gli utenti e gli errori possono essere corretti senza la necessità di apportare modifiche al codice dell'utente (UDF).
  • L'uso delle applicazioni Delta è semplificato per gli utenti, poiché la piattaforma offre flessibilità e resilienza di default e raccoglie numerose metriche dettagliate che possono essere utilizzate per gli avvisi.

Utilizzo in produzione

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

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

Ringraziamenti

Vorremmo ringraziare le seguenti persone che hanno partecipato 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 degli eventi online. Commun. ACM 62(5): 43–49 (2019). DOI: doi.org/10.1145/3312527

Iscriviti al webinar gratuito: «Data Build Tool per il magazzino Amazon Redshift».

Fonte: habr.com

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