Progettare per flessibilità ed efficienza su Dataflow

Questo documento spiega la progettazione per la flessibilità e l'efficienza (DFE), un insieme di best practice architetturali per la creazione di pipeline Dataflow resilienti.

Il passaggio da vincoli rigidi dell'infrastruttura a definizioni flessibili delle risorse ti consente di:

  • Massimizza l'ottenimento delle risorse di computing.
  • Garantire la scalabilità automatica senza interruzioni durante i periodi di forte domanda regionale.
  • Evita ritardi nell'avvio della pipeline ed elimina i colli di bottiglia della capacità.

Ad esempio, anziché limitare la pipeline a un tipo di macchina specifico in una zona, ad esempio richiedendo n1-standard-4 worker in us-central1-a, puoi impostare esigenze minime di risorse (come 4 vCPU e 16 GB di RAM). Se us-central1-a o la serie di macchine N1 riscontrano limiti di capacità temporanei, Dataflow può eseguire automaticamente il provisioning di VM worker compatibili in altre zone e famiglie di macchine (come E2, N2 o N2D). Questa flessibilità contribuisce a garantire che la pipeline venga avviata e scalata senza attendere un singolo pool hardware vincolato.

Questo documento è destinato a data engineer, cloud architect e amministratori di piattaforma che gestiscono i carichi di lavoro Dataflow e vogliono ottimizzare l'affidabilità, la velocità effettiva e la disponibilità dell'infrastruttura della pipeline.

Panoramica di DFE

Dataflow è un servizio di elaborazione dei dati serverless completamente gestito che esegue il provisioning dinamico delle istanze di macchine virtuali (VM) Compute Engine per eseguire le pipeline Apache Beam. Nelle pipeline di elaborazione batch e streaming su larga scala, i pool di worker vengono spesso scalati fino a decine o centinaia di istanze VM.

Le pipeline configurate con vincoli rigidi dell'infrastruttura sono soggette a ritardi nel provisioning durante i periodi di forte domanda. Alcuni esempi di vincoli rigidi sono:

  • Codifica hardcoded di un singolo tipo di macchina, ad esempio n1-standard-4.
  • Bloccare la pipeline in una zona Compute Engine specifica.

Se quel tipo di macchina o zona specifico registra una domanda elevata temporanea, Dataflow non può allocare risorse di calcolo. Ciò può causare ritardi o errori di provisioning come ZONE_RESOURCE_POOL_EXHAUSTED o RESOURCE_POOL_EXHAUSTED.

L'utilizzo dei principi DFE ti aiuta a spostare l'architettura della pipeline da dichiarazioni di infrastruttura rigide e statiche a definizioni di risorse flessibili e basate sui requisiti. Questa flessibilità consente a Dataflow di distribuire dinamicamente il calcolo in diversi pool di hardware disponibili in Google Cloud, aiutandoti a massimizzare l'ottenibilità del calcolo e a ridurre al minimo l'overhead operativo.

Best practice per DFE

Adotta le seguenti best practice per massimizzare l'ottenimento di risorse di calcolo, migliorare la reattività dello scalabilità automatica e creare pipeline resilienti.

Abilita la selezione automatica delle VM

Anziché codificare in modo rigido un tipo di macchina statico con l'opzione pipeline worker machine type, utilizza la selezione automatica della VM con i suggerimenti per le risorse di Apache Beam. Quando specifichi i requisiti minimi di risorse (min_ram o cpu_count), Dataflow attiva automaticamente la flessibilità delle istanze e il provisioning dei worker da un elenco di tipi di macchine compatibili.

Supporto dei workload:

  • Pipeline batch:il dimensionamento corretto e la selezione automatica delle VM vengono attivati automaticamente quando specifichi i suggerimenti per le risorse.
  • Pipeline di streaming:il dimensionamento corretto richiede l'impostazione dell'opzione pipeline --experiments=enable_streaming_rightfitting, insieme alla scalabilità automatica orizzontale (abilitata per impostazione predefinita) e a Streaming Engine (--enable_streaming_engine).

Per configurare la selezione automatica delle VM, specifica i requisiti minimi delle risorse (min_ram o cpu_count) a livello di pipeline utilizzando le opzioni della riga di comando, le opzioni della pipeline SDK o i parametri di esecuzione del modello flessibile. Per istruzioni di configurazione dettagliate ed esempi di codice per Java e Python, vedi Utilizzare i suggerimenti per le risorse.

Utilizza il posizionamento dei worker regionali (evita il pinning zonale)

Configura Dataflow per pianificare dinamicamente le VM worker in qualsiasi zona integra all'interno della regione scelta.

Specifica l'opzione della pipeline --region e ometti --zone e --worker_zone. Ad esempio:

--region=us-central1

Disaccoppia lo stato e mescola utilizzando i servizi gestiti

Le pipeline che non utilizzano servizi di backend gestiti eseguono operazioni di shuffling dei dati e archiviazione dello stato di streaming direttamente sui dischi e nella memoria delle VM worker. Questo stretto accoppiamento richiede dischi di lavoro più grandi e lega la sopravvivenza del workload a istanze VM specifiche, rendendo più difficile la sostituzione dei worker durante i vincoli di capacità.

  • Per i job batch, utilizza Dataflow Shuffle: Dataflow Shuffle è abilitato per impostazione predefinita per le pipeline batch in esecuzione sui tipi di macchine worker supportati e scarica le operazioni di shuffling dalle VM worker a un servizio di backend dedicato e gestito da Google.
  • Per i job di streaming, utilizza Streaming Engine: Streaming Engine esegue l'offload dell'archiviazione dello stato della finestra e della gestione dei timer dalle VM worker a un'infrastruttura di backend specializzata e altamente reattiva. Per le pipeline che utilizzano l'SDK Apache Beam 2.30.0 o versioni successive, Streaming Engine è abilitato per impostazione predefinita. Per abilitarlo in modo esplicito, passa l'opzione della pipeline --enable_streaming_engine.

Utilizza la pianificazione flessibile delle risorse (FlexRS) per le pipeline batch

Per i carichi di lavoro batch non urgenti, come ETL notturni, importazione di data lake o rollup giornalieri, utilizza Flexible Resource Scheduling (FlexRS).

Per attivare FlexRS, imposta l'opzione della pipeline dell'obiettivo flexRS:

  • Per le pipeline Python: --flexrs_goal=COST_OPTIMIZED
  • Per le pipeline Java: --flexRSGoal=COST_OPTIMIZED

Configura i tipi di VM di avvio flessibile per i modelli flessibili

Quando avvii pipeline utilizzando modelli flessibili, la VM di avvio della pipeline utilizza per impostazione predefinita e2-standard-2. La VM predefinita funziona nella maggior parte dei casi, ma se riscontri limitazioni di capacità, puoi personalizzare la configurazione utilizzando l'opzione --launcher-machine-type quando esegui il comando gcloud dataflow flex-template run:

gcloud dataflow flex-template run my-job \
    --template-file-gcs-location="gs://my-bucket/template.json" \
    --region="us-central1" \
    --launcher-machine-type="n2-standard-2"

Considerazioni operative e compromessi

L'adozione delle best practice di DFE migliora significativamente l'ottenibilità del calcolo, la reattività della scalabilità automatica e l'affidabilità operativa. Tuttavia, quando progetti la tua architettura, prendi in considerazione i seguenti fattori operativi e compromessi:

Considerazioni sulla selezione automatica delle VM

  • Affidabilità rispetto alle prestazioni di picco: la selezione automatica delle VM dà la priorità all'affidabilità dell'avvio dei job e all'ottenimento di risorse di calcolo rispetto alle prestazioni di picco dell'esecuzione. Poiché Dataflow esegue il provisioning da più famiglie di macchine candidate (come E2, N2, N4 e N2D), le prestazioni e il throughput di runtime potrebbero variare leggermente a seconda della famiglia di macchine di cui viene eseguito il provisioning. Per i workload a elevato utilizzo di risorse di calcolo con SLA di esecuzione rigorosi, testa la pipeline con la selezione automatica delle VM per stabilire una baseline delle prestazioni prima di eseguirne il deployment su larga scala. Se un workload richiede una piattaforma hardware o una velocità di clock specifica e puoi tollerare i limiti di capacità, puoi continuare a impostare un tipo di macchina specifico.
  • Quota di Compute Engine nelle famiglie candidate:poiché la selezione automatica delle VM può eseguire il provisioning dei worker da più famiglie di macchine candidate, assicurati che il tuo Google Cloud progetto disponga di una quota di vCPU e memoria di Compute Engine sufficiente per ogni famiglia candidata nella regione di destinazione. Se si verifica una carenza di capacità nella famiglia principale e il tuo progetto non dispone di quota per la famiglia di riserva, il provisioning dei worker non va a buon fine e viene visualizzato l'errore QUOTA_EXCEEDED.
  • Prerequisito per la pipeline di streaming:per le pipeline di streaming, la scelta della VM giusta e la selezione automatica della VM non sono abilitate per impostazione predefinita. Devi specificare esplicitamente --experiments=enable_streaming_rightfitting e assicurarti che siano attivi sia Streaming Engine (--enable_streaming_engine) sia la scalabilità automatica orizzontale.
  • Esclusioni di configurazione:la selezione automatica delle VM viene ignorata automaticamente o non è supportata se configuri una delle funzionalità o delle opzioni nella tabella seguente:

    Funzionalità Opzione di configurazione o flag Note
    Tipi di macchina espliciti --worker_machine_type o --machine_type (Python)
    --workerMachineType (Java)
    La selezione automatica delle VM viene ignorata a favore del tipo di macchina specificato.
    Tipi di disco personalizzati, IOPS sottoposte a provisioning o throughput --disk_type, --disk_provisioned_iops o --disk_provisioned_throughput_mibps La selezione automatica delle VM viene ignorata. È supportata l'impostazione di una dimensione del disco personalizzata con --disk_size_gb.
    Piattaforme CPU minime --min_cpu_platform (Python)
    --minCpuPlatform (Java)
    L'impostazione di una piattaforma CPU minima ignora la selezione automatica delle VM.
    Confidential VM --experiments=enable_confidential_compute Le istanze Confidential VM non sono supportate con la selezione automatica delle VM.
    Acceleratori GPU o TPU Suggerimento per la risorsa --dataflow_service_options=worker_accelerator=... o accelerator La selezione automatica delle VM si applica solo ai workload senza acceleratori.
    Dataflow Prime --dataflow_service_options=enable_prime Dataflow Prime utilizza la scalabilità automatica verticale e il dimensionamento dinamico anziché la selezione automatica delle VM.
    Pianificazione flessibile delle risorse (FlexRS) --flexrs_goal=COST_OPTIMIZED (Python)
    --flexRSGoal=COST_OPTIMIZED (Java)
    FlexRS gestisce il proprio pool di worker e il buffer di pianificazione.

Compromessi di Flexible Resource Scheduling (FlexRS)

  • Finestra di ritardo della pianificazione:FlexRS può introdurre un buffer di pianificazione fino a 6 ore prima dell'inizio dell'esecuzione del job. Non utilizzare FlexRS per pipeline con SLA rigidi per il tempo di completamento o dipendenze downstream strette.

Posizionamento regionale e località dei dati

  • Prerequisito del servizio gestito: il posizionamento dei worker regionali è supportato solo per i job che utilizzano Dataflow Shuffle per l'elaborazione batch o Streaming Engine per lo streaming. I job che non utilizzano questi servizi di backend gestiti utilizzano il posizionamento automatico delle zone, che seleziona una singola zona migliore all'interno della regione.
  • Località dei dati e uscita tra regioni:il posizionamento regionale distribuisce i worker tra le zone disponibili all'interno della regione scelta. Per ridurre al minimo la latenza di rete ed evitare addebiti per l'uscita di rete tra regioni, assicurati che tutte le origini e i sink di dati (come i bucket Cloud Storage, i set di dati BigQuery e gli argomenti Pub/Sub) si trovino nella stessa regione del job Dataflow.

Prenotazioni di Compute Engine

  • Affinità di prenotazione:i job Dataflow on demand utilizzano automaticamente le prenotazioni di Compute Engine corrispondenti che utilizzano l'affinità di prenotazione ANY. Tuttavia, la selezione automatica delle VM non supporta l'utilizzo di istanze da prenotazioni specifiche denominate.
  • Idoneità per carichi di lavoro temporanei: le prenotazioni di Compute Engine non sono generalmente consigliate per carichi di lavoro batch a breve durata, con scalabilità automatica o soggetti a picchi di utilizzo. Inoltre, la creazione di nuove prenotazioni durante una carenza zonale attiva non va a buon fine a causa degli stessi vincoli di capacità della creazione di VM on demand.

Passaggi successivi