Utilizza Dataflow Portable Runner

Quando utilizzi Dataflow per eseguire la pipeline, il runner Dataflow carica il codice della pipeline e le dipendenze in un bucket Cloud Storage e crea un job Dataflow. Questo job Dataflow esegue la pipeline su risorse gestite in Google Cloud.

  • Per le pipeline batch che utilizzano le versioni 2.54.0 o successive dell'SDK Apache Beam Java, Portable Runner è abilitato per impostazione predefinita.
  • Per le pipeline che utilizzano l'SDK Apache Beam Java, Portable Runner è necessario quando si eseguono pipeline multilingue, si utilizzano container personalizzati o si utilizzano pipeline di flussi di modifiche Spanner o Bigtable. Per le pipeline di flusso Java su larga scala, utilizza Streaming Java (predefinito) Runner.
  • Per le pipeline che utilizzano le versioni 2.21.0 o successive dell'SDK Apache Beam Python, Portable Runner è abilitato per impostazione predefinita. Per le pipeline che utilizzano le versioni 2.45.0 e successive dell'SDK Apache Beam Python, Dataflow Portable Runner è l'unico runner Dataflow disponibile.
  • Per l'SDK Apache Beam per Go, Portable Runner è l'unico runner Dataflow disponibile.

Portable Runner utilizza un'architettura basata su servizi che avvantaggia molte pipeline:

Consigli sull'utilizzo

A meno che la pipeline non sia interessata da uno dei limiti descritti in questo documento, utilizza le seguenti indicazioni per selezionare il runner appropriato in base ai requisiti della pipeline:

Seleziona Portable Runner per:

  • Tutte le pipeline batch.
  • Pipeline di flusso che utilizzano l'SDK Apache Beam Python o l'SDK Apache Beam per Go.
  • Pipeline di flusso che richiedono il supporto GPU per l'elaborazione di machine learning.
  • Pipeline che utilizzano una delle funzionalità solo di Portable Runner.

Seleziona Streaming Java Runner per:

  • Pipeline di flusso che utilizzano l'SDK Apache Beam Java.

Limitazioni e restrizioni

Dataflow Portable Runner presenta i seguenti requisiti e limitazioni:

  • Funzionalità solo di Portable Runner: le seguenti funzionalità sono supportate solo su Portable Runner:
    • Trasformazioni gestite e RunInference in Java.
    • Container personalizzati.
    • VM worker basate su ARM.
    • DoFn divisibili. È in fase di esplorazione il supporto per Non-Portable Runner.
  • Dataflow Portable Runner richiede Streaming Engine per i job di flusso.
  • Poiché Dataflow Portable Runner richiede Streaming Engine per i job di flusso, qualsiasi trasformazione Apache Beam che richiede Dataflow Portable Runner richiede anche l'utilizzo di Streaming Engine per i job di flusso. Ad esempio, il connettore I/O Pub/Sub Liteper l'SDK Apache Beam per Python è una trasformazione cross-language che richiede Dataflow Portable Runner. Se provi a disattivare Streaming Engine per un job o un modello che utilizza questa trasformazione, il job non riesce.
  • Per le pipeline di flusso che utilizzano l'SDK Apache Beam Java, le classi MapState e SetState non sono supportate con Portable Runner. Per utilizzare le classi MapState e SetState con le pipeline Java, attiva Streaming Engine, disattiva Portable Runner e utilizza l'SDK Apache Beam versione 2.58.0 o successive.
  • Per le pipeline batch e di flusso che utilizzano l'SDK Apache Beam Java, la classe AfterSynchronizedProcessingTime non è supportata.
  • Sebbene Portable Runner sia più scalabile di Non-Portable Runner in molti casi, la memoria utilizzata potrebbe essere maggiore per lo sharding fisso.
  • I modelli classici di Dataflow non possono essere eseguiti con una versione del runner Dataflow diversa da quella con cui sono stati creati. Ciò significa che i modelli classici forniti da Google non possono attivare Portable Runner. Per attivare Portable Runner per i modelli personalizzati, imposta il flag appropriato quando crei il modello. Questo flag dipende dalla versione dell'SDK che stai utilizzando:
    • Versioni 2.74 e successive dell'SDK Beam: --experiments=enable_portable_runner.
    • Versioni 2.73 e precedenti dell'SDK Beam: --experiments=use_runner_v2.
  • A causa di un problema di scalabilità automatica noto, Portable Runner è disattivato per impostazione predefinita per le pipeline batch Java che richiedono l'elaborazione con stato. Puoi comunque attivare Portable Runner per queste pipeline (vedi Attivare Portable Runner), ma le prestazioni della pipeline potrebbero essere gravemente limitate.

  • In alcune pipeline, Portable Runner può aumentare la frequenza degli errori di coerenza. Nei file di log potrebbe essere visualizzato il seguente errore: "Internal consistency check failed, the output is likely incorrect. Please retry the job". Una possibile mitigazione consiste nell'aggiungere una trasformazione Reshuffle dopo il passaggio Join/GroupByKey. Se il tasso di errore non è tollerabile e la mitigazione non risolve il problema, prova a disattivare Portable Runner.

Attivare Portable Runner

Per attivare Dataflow Portable Runner, segui le istruzioni di configurazione per l'SDK Apache Beam.

Java

Dataflow Portable Runner richiede le versioni 2.30.0 o successive dell'SDK Apache Beam Java, con la versione 2.44.0 o successive consigliata.

Per le pipeline batch che utilizzano le versioni 2.54.0 o successive dell'SDK Apache Beam Java, Portable Runner è abilitato per impostazione predefinita.

Per attivare Portable Runner, esegui il job con il valore dell'esperimento corrispondente alla versione dell'SDK Beam e al tipo di pipeline:

  • Versioni 2.74 e successive dell'SDK Beam: enable_portable_runner.
  • Versioni 2.73 e precedenti dell'SDK Beam: use_runner_v2.

Per saperne di più, consulta Impostare le opzioni sperimentali della pipeline.

Python

Per le pipeline che utilizzano le versioni 2.21.0 o successive dell'SDK Apache Beam Python, Portable Runner è abilitato per impostazione predefinita.

Dataflow Portable Runner non è supportato con le versioni 2.20.0 e precedenti dell'SDK Apache Beam Python.

In alcuni casi, la pipeline potrebbe non utilizzare Portable Runner anche se viene eseguita su una versione dell'SDK supportata. Per eseguire il job con Portable Runner, imposta il valore dell'esperimento corrispondente alla versione dell'SDK Beam:

  • Versioni 2.74 e successive dell'SDK Beam: enable_portable_runner.
  • Versioni 2.73 e precedenti dell'SDK Beam: use_runner_v2.

Per saperne di più, consulta Impostare le opzioni sperimentali della pipeline.

Vai

Dataflow Portable Runner è l'unico runner Dataflow disponibile per l'SDK Apache Beam per Go. Portable Runner è abilitato per impostazione predefinita.

Disattivare Portable Runner

Per disattivare Dataflow Portable Runner, segui le istruzioni di configurazione per l'SDK Apache Beam.

Java

Per disattivare Portable Runner, imposta il valore dell'esperimento corrispondente alla versione dell'SDK Beam:

  • Versioni 2.74 e successive dell'SDK Beam:
    • Per il batch, utilizza disable_portable_runner.
    • Per il flusso, utilizza enable_streaming_java_runner.
  • Versioni 2.73 e precedenti dell'SDK Beam: disable_runner_v2.

Il job verrà impostato per impostazione predefinita su Non-Portable Runner o Streaming Java Runner per il flusso. Per saperne di più, consulta Impostare le opzioni sperimentali della pipeline.

Python

La disattivazione di Portable Runner non è supportata con le versioni 2.45.0 e successive dell'SDK Apache Beam Python.

Per le versioni precedenti dell'SDK Python, se il job viene identificato come utilizzo dell'esperimento auto_runner_v2, puoi disattivare Portable Runner impostando l'esperimento disable_runner_v2. Per saperne di più, consulta Impostare le opzioni sperimentali della pipeline.

Vai

Dataflow Portable Runner non può essere disattivato in Go. Portable Runner è l'unico runner Dataflow disponibile per l'SDK Apache Beam per Go.

Monitorare il job

Utilizza l'interfaccia di monitoraggio per visualizzare le metriche dei job Dataflow, come l'utilizzo della memoria, l'utilizzo della CPU e altro ancora.

I log delle VM worker sono disponibili tramite Esplora log e l' interfaccia di monitoraggio Dataflow. I log delle VM worker includono i log del processo di harness del runner e i log dei processi dell'SDK. Puoi utilizzare i log delle VM per risolvere i problemi relativi al job.

Risolvere i problemi di Portable Runner

Per risolvere i problemi relativi ai job che utilizzano Dataflow Portable Runner, segui i passaggi standard per la risoluzione dei problemi relativi alle pipeline. Il seguente elenco fornisce ulteriori informazioni sul funzionamento di Dataflow Portable Runner:

  • I job Dataflow Portable Runner eseguono due tipi di processi sulla VM worker: il processo SDK e il processo di harness del runner. A seconda della pipeline e del tipo di VM, potrebbero essere presenti uno o più processi SDK, ma è presente un solo processo di harness del runner per VM.
  • I processi SDK eseguono il codice utente e altre funzioni specifiche per la lingua. Il processo di harness del runner gestisce tutto il resto.
  • Il processo di harness del runner attende che tutti i processi SDK si connettano prima di iniziare a richiedere il lavoro a Dataflow.
  • I job potrebbero essere ritardati se la VM worker scarica e installa le dipendenze durante l'avvio del processo SDK. Se si verificano problemi durante un processo SDK, ad esempio durante l'avvio o l'installazione delle librerie, il worker segnala il suo stato come non integro. Se i tempi di avvio aumentano, attiva l'API Cloud Build nel tuo progetto e invia la pipeline con il seguente parametro: --prebuild_sdk_container_engine=cloud_build.
  • Poiché Dataflow Portable Runner utilizza il checkpointing, ogni worker potrebbe attendere fino a cinque secondi durante il buffering delle modifiche prima di inviarle per l'ulteriore elaborazione. Di conseguenza, è prevista una latenza di circa sei secondi.
  • Per diagnosticare i problemi nel codice utente, esamina i log dei worker dai processi SDK. Se riscontri errori nei log di harness del runner, contatta l'assistenza per segnalare un bug.
  • Per eseguire il debug degli errori comuni relativi alle pipeline multilingue di Dataflow, consulta la guida Suggerimenti per le pipeline multilingue.