Dataflow è progettato per eseguire pipeline di elaborazione dei dati di grandi dimensioni distribuendo il lavoro in un pool gestito di istanze di calcolo. Comprendere come Dataflow parallelizza l'elaborazione ti aiuta a progettare pipeline efficienti, evitare colli di bottiglia delle prestazioni e ottimizzare i costi delle risorse.
Questa pagina spiega come Dataflow parallelizza l'elaborazione dei dati, come gestisce e scala l'esecuzione, i fattori comuni che limitano il parallelismo e le tecniche che puoi utilizzare per ottimizzare la velocità effettiva della pipeline.
Modelli di parallelismo: orizzontale e verticale
Dataflow raggiunge il parallelismo utilizzando due strategie complementari:
Parallelismo orizzontale: i dati della pipeline vengono partizionati ed elaborati su più istanze worker (macchine virtuali) contemporaneamente. Dataflow può regolare automaticamente le dimensioni del pool di worker in base alla domanda del carico di lavoro tramite la scalabilità automatica orizzontale. Per impostazione predefinita, Dataflow imposta un limite di risorse di 4000 worker per job, che può essere modificato utilizzando le richieste di quota.
Parallelismo verticale: più core e thread della CPU all'interno di una singola istanza worker elaborano i dati della pipeline contemporaneamente. Ogni VM worker esegue processi worker e thread di harness per utilizzare le risorse di calcolo disponibili. Con lo scaling dinamico dei thread, Dataflow può regolare il numero di thread attivi per worker nelle pipeline batch in base all'utilizzo della CPU e alla memoria disponibile. In Dataflow Prime, la scalabilità automatica verticale ridimensiona dinamicamente la memoria e il calcolo allocati ai worker.
Unità di lavoro e gerarchia di esecuzione
Per distribuire l'elaborazione tra worker e thread, Dataflow divide le pipeline Apache Beam in unità di lavoro discrete:
- PCollection e partizioni: un
PCollectionrappresenta un set di dati distribuito. Per i dati delimitati (pipeline batch), Dataflow divide il set di dati in suddivisioni o shard. Per i dati illimitati (pipeline di streaming), i dati arrivano continuamente e vengono inseriti come messaggi o partizioni di stream. - Bundle:Dataflow raggruppa gli elementi in bundle arbitrari
per l'elaborazione da parte di un
DoFn. Un bundle è l'unità di errore e di ripetizione: se l'elaborazione di un elemento genera un'eccezione non gestita, viene eseguito un nuovo tentativo per l'intero bundle. Le operazioni con un elevato consumo di memoria possono aumentare la pressione della memoria del worker e causare errori di esaurimento della memoria. - Unione di fasi e passaggi: durante l'ottimizzazione del grafico, Dataflow combina le trasformazioni adiacenti in fasi di esecuzione unite per eliminare l'overhead della materializzazione dei dati intermedi. All'interno di una fase unita, gli elementi vengono elaborati in un ciclo di esecuzione stretto su un singolo thread prima di essere passati alla fase successiva o al limite di rimescolamento.
Per ulteriori dettagli sulla traduzione delle pipeline e sulla generazione dei grafici, consulta Ciclo di vita della pipeline.
Parallelismo gestito e scalabilità automatica
Per impostazione predefinita, Dataflow gestisce automaticamente il parallelismo della pipeline senza richiedere la regolazione manuale delle partizioni nei seguenti modi:
- Scalabilità automatica orizzontale:
- Pipeline batch:valuta il lavoro rimanente totale stimato, il backlog dell'origine e l'utilizzo della CPU per scalare il worker pool verso l'alto o verso il basso per completare il job in modo rapido ed economico.
- Pipeline di streaming:analizza la latenza del sistema, le dimensioni del backlog e l'utilizzo della CPU per aumentare il numero di worker durante i picchi di throughput e fare lo scale down durante i periodi di basso traffico. Per i dettagli, vedi Ottimizzare la scalabilità automatica orizzontale di streaming.
- Bilanciamento dinamico del lavoro (DWR): nelle pipeline batch, Dataflow monitora l'avanzamento delle attività dei singoli worker. Se un worker termina in anticipo o un altro worker è in ritardo a causa della distorsione dei dati (ritardatari), Dataflow divide dinamicamente il lavoro residuo non elaborato dal worker lento e lo riassegna a un worker inattivo. Per maggiori informazioni, consulta Ribilanciamento dinamico del lavoro.
- Scalabilità dinamica dei thread:nelle pipeline batch che utilizzano Portable Runner, regola automaticamente il numero di thread di elaborazione simultanei per worker in base all'utilizzo della CPU e alla memoria disponibile. Per saperne di più, consulta Scalabilità dinamica dei thread.
- Scalabilità automatica verticale: in Dataflow Prime, Dataflow scala dinamicamente la memoria e le risorse di calcolo dei worker per evitare errori di memoria insufficiente e ottimizzare l'utilizzo delle risorse. Per saperne di più, consulta Scalabilità automatica verticale.
Fattori che limitano il parallelismo
Una pipeline potrebbe non raggiungere il parallelismo previsto a causa delle seguenti caratteristiche dei dati o della progettazione del grafico della pipeline:
Origini di input non divisibili
Se una sorgente di ingresso non può essere suddivisa in intervalli indipendenti, Dataflow è costretto a leggere la sorgente in sequenza con un singolo thread di lavoro:
- Compressione di file non divisibili:i formati come
.gz(gzip) o.bzip2(senza indicizzazione) non possono essere letti in parallelo da offset di byte arbitrari. La lettura di un singolo file compresso di grandi dimensioni limita la fase di importazione a un singolo thread finché i dati non vengono decompressi e ridistribuiti. - Soluzione:archivia i dati in formati di file divisibili (come Parquet, Avro o formati compressi con Snappy) o dividi i dati di input in più file più piccoli in Cloud Storage.
Step fusion e fan-out elevato
La fusione dei passaggi migliora il rendimento riducendo l'overhead di serializzazione, ma può limitare inavvertitamente il parallelismo e aumentare la pressione sulla memoria quando un passaggio con parallelismo basso produce un numero elevato di elementi di output (un'operazione di "fan-out" elevato):
- Esempio: un'origine legge cinque file e viene unita a una trasformazione
FlatMapche produce 1.000.000 di elementi di output. Se la trasformazioneFlatMapviene unita alle trasformazioni downstream, tutti i 1.000.000 di elementi continuano a essere eseguiti su un massimo di cinque thread di lavoro, limitando gravemente la velocità effettiva downstream. Inoltre, se le trasformazioni intermedie si espandono in modo significativo in memoria prima del commit, i bundle di grandi dimensioni possono esaurire la memoria di lavoro disponibile. - Risoluzione: inserisci una trasformazione
Redistribute(oReshuffleclassica) tra il passaggio con fan-out elevato e le trasformazioni downstream per interrompere la fusione e ridistribuire il lavoro nel pool di worker. Per il debug dei problemi di memoria correlati, vedi Risolvere gli errori di out of memory.
Inclinazione dei tasti e tasti di scelta rapida
Le operazioni di aggregazione (GroupByKey, CoGroupByKey, Combine.PerKey) raggruppano
gli elementi in base alla chiave associata.
- Collo di bottiglia della hot key:Dataflow indirizza tutti gli elementi con la stessa chiave a un singolo thread di lavoro per l'aggregazione. Se una singola chiave contiene
una percentuale elevata del set di dati totale, il worker diventa un ritardatario
e i worker upstream potrebbero subire una contropressione. Ad esempio, una chiave
nullpredefinita o una chiave di categoria estremamente popolare. - Risoluzione:
- Utilizza Combiner
(
CombineFnoCombine.PerKey) anzichéGroupByKey, se possibile, per consentire a Dataflow di eseguire combinazioni locali parziali prima del rimescolamento. - Aggiungi un prefisso o un suffisso intero casuale ai tasti di scelta rapida (salatura delle chiavi) per distribuire lo spazio delle chiavi tra i worker, seguito da un'aggregazione di secondo livello per unire i risultati salati.
- Utilizza Combiner
(
Limitazione del sink downstream
Quando scrivi l'output della pipeline in servizi esterni come database o API di terze parti, un parallelismo elevato può saturare il sistema di destinazione:
- Limitazione: centinaia di thread di lavoro che eseguono chiamate di scrittura simultanee possono causare errori di limite di frequenza, timeout di connessione o degrado del database.
- Risoluzione:
- Limita il parallelismo di scrittura raggruppando gli elementi con
GroupByKeyo utilizzando sink di batch con parallelismo controllato. - Implementa il backoff esponenziale lato client e la logica di ripetizione nei sink
DoFn.
- Limita il parallelismo di scrittura raggruppando gli elementi con
Strategie di ottimizzazione
Per ottimizzare il parallelismo nei job Dataflow, valuta i seguenti approcci:
- Evita fusioni indesiderate con
Redistribute:Redistribute.arbitrarily(): interrompe la fusione dei passaggi e ridistribuisce gli elementi in modo uniforme tra tutti i lavoratori disponibili.Redistribute.byKey(): ribilancia le coppie chiave-valore tra i thread di lavoro preservando la località delle chiavi.- Per esempi di implementazione, consulta Impedire la fusione.
- Monitora i ritardatari e i colli di bottiglia: utilizza i dettagli di esecuzione della console Google Cloud per identificare le fasi con un numero elevato di ritardatari o un avanzamento bloccato:
Passaggi successivi
- Scopri di più sul ciclo di vita della pipeline.
- Esplora la scalabilità automatica orizzontale.
- Scopri di più sul ribilanciamento dinamico del lavoro.
- Consulta le best practice per le pipeline Dataflow.
- Scopri come risolvere i problemi relativi agli errori out of memory.