Netflix, come uno dei principali produttori di video streaming, ha bisogno di un sistema estremamente efficiente per gestire grandi quantità di dati temporali, soprattutto per raccogliere dati sull’uso e comportamento degli utenti. Per questo motivo, il team ha collaborato per risolvere un problema molto comune in Apache Cassandra: la latenza elevata causata da partizioni molto grandi. La strategia sviluppata utilizza una tecnica chiamata "Dynamic Partitioning", mirata a dividere queste partizioni in maniera trasparente e automatica, migliorando così la velocità delle letture.

Che cos’è Dynamic Repartitioning?

Il TimeSeries Abstraction di Netflix serve milioni di richieste al secondo per leggere dati temporali con latenza al millisecondo. Lui fa affidamento su Apache Cassandra perché offre elevate prestazioni in termini di latenza, throughput, costo e maturità operativa. Il problema sorge quando le partizioni si allargano troppo, causando problemi di latenza e possibili tempi di attesa in lettura. Dynamic Repartitioning offre una soluzione trasparente, dividendole senza alterare le applicazioni, e riuscendo a mantenere il partizionamento logico coerente.

Perché le partizioni grandi influenzano negativamente le letture?

In genere, la latenza di lettura sta in un singolo cifra di millisecondi. Tuttavia, quando un certo numero di eventi si accumula in una partizione, essa può diventare estremamente grande ("wide"). Questo fenomeno incrementa la latenza fino a secondi, causando timeout. In casi estremi, i nodi di Cassandra possono presentare pause nella Garbage Collection, utilizzo della CPU elevato e un aumento delle code di thread.

Il TimeSeries gestisce un throughput elevatissimo di letture, il che complica ulteriormente il problema. Sebbene si possa sempre aumentare la capacità con un cluster più grande, l’approccio Netflix ha cercato di migliorare il livello di efficienza senza bisogno di nuove risorse dedicate.

Chi si trova di fronte a partizioni troppo grandi?

Il TimeSeries usa una logica di partizionamento basata su finestre di tempo: TimeSlice, time bucket e event bucket. Permette a Netflix di eseguire query e di eliminare i dati in base al tempo senza creare tombstone entries. Quando un dataset viene creato, l’utente specifica i parametri di utilizzo previsto. Un processo di provisioning effettua simulazioni Monte Carlo per determinare la configurazione ottimale.

Tuttavia, esistono tre casi in cui questa strategia non è sufficiente. Un caso è quando i carichi di lavoro non sono noti o sono stimati male inizialmente. Un altro riguarda quando i carichi di lavoro cambiano nel tempo a causa di nuove richieste o traffico variabile. Infine, ci sono i casi di dati anomali in cui pochi ID ricevono un numero notevolmente elevato di eventi.

Soluzione 1: Re-Partitioning a Tempo

Cassandra fornisce API di introspection come nodetool tablehistograms per analizzare le percentuali di grandezza delle partizioni. Un worker di background monitora questi dati e li presenta tramite una tabella virtuale. Quando le partizioni superano un livello target (frequente tra 2 MB e 10 MB a seconda del carico), il worker calcola un fattore correttivo.

Ad esempio, inizialmente si usavano intervalli di 60 secondi che producevano partizioni piccole (under 10 KB), causando alta lettura di dati (amplificazione) e coda di thread nel sistema. Il DynamicTimeSliceConfigWorker ha proposto cambiamenti per aumentare il tempo di bucket:

    • Ns: mydataset1
    • Osservato: Le partizioni di TimeSlices avevano p99 sotto il target di 10 MB.
    • Proposto: Cambiato il time_bucket da 60s a 604800s.

Questo ha ridotto le latenze e i timeout causati da code di thread. Tuttavia, la soluzione funziona meglio quando la maggior parte del dataset richiede il partizionamento. Non aiuta quando solo una percentuale di ID ha partizioni molto ampie.

Soluzione 2: Partitioning Dinamico per ID

Il Dynamic Partitioning è un processo asincrono che divide le partizioni molto grandi per ID specifico. Funziona al livello ID e ha tre fasi principali: Detection, Planning & Splitting, e Serving Reads.

Detection

Durante le letture, l’applicazione tiene traccia del numero di byte letti per partizione. Quando superano una soglia configurata, il server emette un evento su Kafka:

{

"timeslice": "data20260328",

"timeseriesid": "profileId:123",

"time_bucket": 7,

"event_bucket": 2,

"immutable": true,

"version": "0"

}

La chiave "immutable" identifica una partizione su cui non vengono ricevuti nuovi dati. L’attributo "version" sarà usato in futuro per invalidare modifiche. Il controllo avviene su letture, non su scritture, poiché molte partizioni non richiedono mai un partizionamento. L’implementazione iniziale si concentra su partizioni immutabili per ridurre la complessità.

Planning

La fase di pianificazione legge l’intera partizione una volta per creare un piano di divisione preciso. L’uso di checkpoint permette di riprendere da dove si era interrotto, in caso di fallimento. La tabella wide_row mantiene informazioni sullo stato, i checkpoint e il routing.

Splitting

Il splitting viene passato a una strategia, ad esempio EventBucketPartitionSplitStrategy. Questa assegna ulteriori bucket di evento al bucket di tempo. Per partizioni estremamente ampie, il numero di bucket è limitato per gestire l’amplificazione del caricamento. La divisione distribuisce i dati tra le repliche di Cassandra.

Visione del percorso di lettura

I server TimeSeries caricano periodicamente in cache le chiavi di partizione divise in un Bloom Filter in memoria. Ogni lettura controlla velocemente il Bloom Filter, che risponde con tempi nell’ordine di singole cifre di microsecondi.

{

"presplitdata": {

"timeslice": "data20260328",

"timeseriesid": "6313825",

"time_bucket": 0,

"event_bucket": 2

},

"postsplitdata": {

"timeslice": "widedata202603280",

"eventbucketpartition_strategy": {

"targeteventbuckets": 2,

"starteventbucket": 32

}

}

}

Queste informazioni sono supportate da una cache ad accesso veloce. Il PartitionReader serve direttamente i dati dal nuovo set di partizioni divise. Il vecchio set non viene eliminato, offrendo una via sicura in caso di errori.

Netflix ha eseguito in modo offline la verifica dei dati splittati con i job Spark di Data Bridge. Un’analisi in tempo reale confrontava i dati letti attraverso il percorso precedente e successivo all’aggiornamento.

Confronto tra le due soluzioni

Dimensione Temporizzazione (Soluzione 1) Partizionamento Din
Read original article →
← Back to news