Informazioni sul parallelismo in Dataflow

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 PCollection rappresenta 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 FlatMap che produce 1.000.000 di elementi di output. Se la trasformazione FlatMap viene 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 (o Reshuffle classica) 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 null predefinita o una chiave di categoria estremamente popolare.
  • Risoluzione:
    1. Utilizza Combiner (CombineFn o Combine.PerKey) anziché GroupByKey, se possibile, per consentire a Dataflow di eseguire combinazioni locali parziali prima del rimescolamento.
    2. 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.

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 GroupByKey o utilizzando sink di batch con parallelismo controllato.
    • Implementa il backoff esponenziale lato client e la logica di ripetizione nei sink DoFn.

Strategie di ottimizzazione

Per ottimizzare il parallelismo nei job Dataflow, valuta i seguenti approcci:

Passaggi successivi