Informazioni sul passaggio del job Dataflow

Nell'interfaccia di monitoraggio di Dataflow, il riquadro Informazioni sul passaggio mostra informazioni sui singoli passaggi di un job. Un passaggio rappresenta una singola trasformazione nella pipeline. Le trasformazioni composite contengono passaggi secondari.

Il riquadro Informazioni sul passaggio mostra le seguenti informazioni:

  • Metriche per il passaggio.
  • Informazioni sulle raccolte di input e output del passaggio.
  • Quali fasi corrispondono a questo passaggio.
  • Metriche di input secondario

Utilizza il riquadro Informazioni sul passaggio per capire il rendimento del job in ogni passaggio e per trovare i passaggi che possono essere potenzialmente ottimizzati.

Visualizzare le informazioni sul passaggio

Per visualizzare le informazioni sul passaggio:

  1. Nella Google Cloud console, vai alla pagina Dataflow > Job.

    Vai a Job

  2. Seleziona un job.

  3. Fai clic sulla scheda Grafico del job per visualizzare il grafico del job. Il grafico del job rappresenta ogni passaggio della pipeline come una casella.

  4. Fai clic su un passaggio. Le informazioni sul passaggio vengono visualizzate nel riquadro Informazioni sul passaggio.

  5. Per visualizzare i passaggi secondari di una trasformazione composita, fai clic sulla Espandi nodo freccia.

Metriche dei passaggi

Il riquadro Informazioni sul passaggio mostra le seguenti metriche per il passaggio.

Watermark e ritardo del sistema

Il watermark del sistema è il timestamp più recente per il quale sono stati elaborati completamente tutti i tempi degli eventi. Il ritardo del watermark del sistema è il tempo massimo di attesa per l'elaborazione di un elemento di dati.

Watermark e ritardo dei dati

Il watermark dei dati è il timestamp che indica il tempo di completamento stimato dell'input di dati per questo passaggio. Il ritardo del watermark dei dati è la differenza tra il tempo dell'ultimo evento di input e il watermark dei dati.

Tempo totale di esecuzione

Il tempo totale di esecuzione è il tempo approssimativo totale trascorso in tutti i thread di tutti i worker per le seguenti azioni:

  • Inizializzazione del passaggio
  • Elaborazione dei dati
  • Data shuffling
  • Fine del passaggio

Per i passaggi compositi, il tempo totale di esecuzione è uguale alla somma del tempo trascorso nei passaggi dei componenti.

Il tempo totale di esecuzione può aiutarti a identificare i passaggi lenti e a diagnosticare quale parte della pipeline richiede più tempo del necessario.

Stato collo di bottiglia

Se Dataflow rileva un collo di bottiglia, viene mostrato un avviso, insieme alla causa, se nota. Per saperne di più, vedi Risolvere i colli di bottiglia.

Latenza massima operazione

La latenza massima dell'operazione è il tempo massimo impiegato in questo passaggio per elaborare i messaggi in entrata o le scadenze delle finestre. Questa metrica viene misurata in modo aggregato tra i passaggi raggruppati in una singola fase, perciò il valore rappresenta l'intera fase.

Parallelismo delle chiavi

Il parallelismo delle chiavi è il numero approssimativo di chiavi in uso per l'elaborazione dei dati in questo passaggio.

Informazioni sul windowing

Per i passaggi di aggregazione nelle pipeline di streaming, il riquadro Informazioni sul passaggio mostra le informazioni sul windowing derivate dalla strategia di windowing della trasformazione ParDo della pipeline. Ad esempio:

  • Funzione di windowing: org.apache.beam.sdk.transforms.windowing.SlidingWindows
  • Periodo della finestra: 30 sec
  • Dimensione della finestra: 5 min
  • Offset di inizio della finestra: 0 ms
  • Modalità di accumulo del windowing: DISCARDING
  • Fattore di amplificazione della scrittura del windowing: 10

I seguenti campi di windowing forniscono un contesto importante per il rendimento dello streaming e il volume di dati:

  • Modalità di accumulo del windowing: indica se lo stato della finestra viene mantenuto tra le attivazioni dei trigger. Per saperne di più, vedi Modalità di accumulo delle finestre nella documentazione di Apache Beam.
  • Fattore di amplificazione della scrittura del windowing: il numero medio di finestre a cui viene assegnato ogni elemento in entrata, calcolato comeWindow Size / Window Period. Ad esempio, una dimensione della finestra di 5 minuti (300 sec) con un periodo di scorrimento di 30 secondi (30 sec) assegna ogni elemento a 10 finestre sovrapposte (300 / 30 = 10), con un fattore di amplificazione della scrittura di 10. Un fattore di amplificazione maggiore di 1 moltiplica il volume degli elementi downstream, il traffico di shuffle e l'archiviazione dello stato.

Raccolte di input/output

Il riquadro Informazioni sul passaggio mostra le seguenti informazioni su ciascuna delle raccolte di input e output del passaggio:

  • Grafico del throughput. Questo grafico mostra il throughput della raccolta. Puoi visualizzare il grafico come elementi al secondo o come byte al secondo. Per saperne di più su questa metrica, vedi Throughput.

  • Conteggio degli elementi aggiunti alla raccolta.

  • Dimensione stimata della raccolta, in byte.

Fasi ottimizzate

Una fase rappresenta una singola unità di lavoro eseguita da Dataflow. Quando selezioni un passaggio nel grafico del job, il riquadro Informazioni sul passaggio mostra i nomi delle fasi che eseguono questo passaggio, insieme allo stato attuale, ad esempio in esecuzione, arrestato o riuscito.

Per visualizzare ulteriori informazioni sulle fasi del job, utilizza la scheda Dettagli di esecuzione.

Metriche di input secondario

Un input secondario è un input aggiuntivo a cui una trasformazione può accedere ogni volta che elabora un elemento. Se una trasformazione crea o utilizza un input secondario, il riquadro Informazioni secondarie mostra le metriche per la raccolta di input secondario.

Se una trasformazione composita crea o utilizza un input secondario, espandi la trasformazione composita finché non vedi la trasformazione secondaria specifica che crea o utilizza l'input secondario. Seleziona la trasformazione secondaria per visualizzare le metriche di input secondario.

Trasformazioni che creano un input secondario

Se una trasformazione crea una raccolta di input secondario, la sezione Metriche di input secondario mostra il nome della raccolta, insieme alle seguenti metriche:

  • Tempo trascorso in scrittura: il tempo trascorso a scrivere la raccolta di input secondario.
  • Byte scritti: il numero totale di byte scritti nella raccolta di input secondario.
  • Ora e byte letti dall'input secondario: una tabella che contiene metriche aggiuntive per tutte le trasformazioni che utilizzano la raccolta di input secondario, chiamate consumer di input secondario.

La tabella Ora e byte letti dall'input secondario contiene le seguenti informazioni per ogni consumer di input secondario:

  • Consumer di input secondario: il nome della trasformazione del consumer di input secondario.
  • Tempo trascorso a leggere: il tempo trascorso da questo consumer a leggere la raccolta di input secondario.
  • Byte letti: il numero di byte letti da questo consumer dalla raccolta di input secondario.

L'immagine seguente mostra le metriche di input secondario per una trasformazione che crea una raccolta di input secondario:

Metriche di input secondario mostrate nel riquadro Informazioni passaggio

Il grafico del job ha una trasformazione composita espansa (MakeMapView). La trasformazione secondaria che crea l'input secondario (CreateDataflowView) è selezionata e le metriche di input secondario sono visibili nel riquadro Informazioni sul passaggio.

Trasformazioni che utilizzano input secondari

Se una trasformazione utilizza uno o più input secondari, la sezione Metriche di input secondario mostra la tabella Ora e byte letti dall'input secondario. Questa tabella contiene le seguenti informazioni per ogni raccolta di input secondario:

  • Raccolta di input secondario: il nome della raccolta di input secondario.
  • Tempo trascorso a leggere: il tempo trascorso dalla trasformazione a leggere questa raccolta di input secondario.
  • Byte letti: il numero di byte letti dalla trasformazione da questa raccolta di input secondario.

L'immagine seguente mostra le metriche di input secondario per una trasformazione che legge da una raccolta di input secondario.

Metriche di input secondario mostrate nel riquadro Informazioni passaggio

La trasformazione JoinBothCollections legge da una raccolta di input secondario. JoinBothCollections è selezionata nel grafico del job e le metriche di input secondario sono visibili nel riquadro Informazioni sul passaggio.

Identificare i problemi di rendimento dell'input secondario

Gli input secondari possono influire sul rendimento della pipeline. Quando la pipeline utilizza un input secondario, Dataflow scrive la raccolta in un livello persistente, ad esempio un disco, e le trasformazioni leggono da questa raccolta persistente. Queste operazioni di lettura e scrittura influiscono sul tempo di esecuzione del job.

La reiterazione è un problema comune di rendimento dell'input secondario. Se il PCollection di input secondario è troppo grande, i worker non possono memorizzare nella cache l'intera raccolta in memoria. Di conseguenza, i worker devono leggere ripetutamente dalla raccolta di input secondario persistente.

Nell'immagine seguente, le metriche di input secondario mostrano che il numero totale di byte letti dalla raccolta di input secondario è molto maggiore della dimensione della raccolta, indicata come byte totali scritti. La raccolta di input secondario è di 563 MB e la somma dei byte letti dalle trasformazioni che la utilizzano è di quasi 12 GB.

Esempio di reiterazione

Per migliorare il rendimento di questa pipeline, riprogetta l'algoritmo in modo da evitare di iterare o recuperare i dati di input secondario. In questo esempio, la pipeline crea il prodotto cartesiano di due raccolte. L'algoritmo scorre l'intera raccolta di input secondario per ogni elemento della raccolta principale. Puoi migliorare il pattern di accesso della pipeline raggruppando più elementi della raccolta principale. Questa modifica riduce il numero di volte in cui i worker devono rileggere la raccolta di input secondario.

Un altro problema di prestazioni comune può verificarsi se la pipeline esegue un join applicando un ParDo con uno o più input secondari di grandi dimensioni. In questo caso, i worker trascorrono una percentuale elevata del tempo di elaborazione per l'operazione di join leggendo dalle raccolte di input secondario.

L'immagine seguente mostra le metriche di input secondario per questo problema:

Un esempio di join di input laterale costoso

La trasformazione JoinBothCollections ha un tempo di elaborazione totale di oltre 18 minuti. I worker trascorrono la maggior parte del tempo di elaborazione (10 minuti) leggendo dalla raccolta di input secondario di 10 GB. Per migliorare il rendimento di questa pipeline, utilizza CoGroupByKey anziché gli input secondari.