Questa pagina descrive come eseguire un job Dataflow utilizzando un modello flessibile. I modelli flessibili consentono di pacchettizzare una pipeline Dataflow in modo da poterla eseguire senza disporre di un ambiente di sviluppo Apache Beam.
Autorizzazioni obbligatorie
Quando esegui un modello flessibile, Dataflow crea un job per te. Per creare il job, il account di servizio Dataflow deve disporre della seguente autorizzazione:
dataflow.serviceAgent
Quando utilizzi Dataflow per la prima volta, il servizio assegna questo ruolo, quindi non devi concedere questa autorizzazione.
Per impostazione predefinita, il account di servizio Compute Engine viene utilizzato per le VM di avvio e le VM worker. Il account di servizio deve disporre dei seguenti ruoli e funzionalità:
- Storage Object Admin (
roles/storage.objectAdmin) - Visualizzatore (
roles/viewer) - Dataflow Worker (
roles/dataflow.worker) - Accesso in lettura e scrittura al bucket di gestione temporanea
- Accesso in lettura all'immagine del modello flessibile
Per concedere l'accesso in lettura e scrittura al bucket di gestione temporanea, puoi utilizzare il ruolo
Storage Object Admin (roles/storage.objectAdmin). Per saperne di più,
consulta Ruoli IAM per Cloud Storage.
Per concedere l'accesso in lettura all'immagine del modello flessibile, puoi utilizzare il ruolo
Storage Object Viewer (roles/storage.objectViewer). Per saperne di più,
consulta Configurare il controllo dell'accesso.
Eseguire un modello flessibile
Per eseguire un modello flessibile, utilizza il
gcloud dataflow flex-template run
comando:
gcloud dataflow flex-template run JOB_ID \ --template-file-gcs-location gs://TEMPLATE_FILE_LOCATION \ --region REGION \ --staging-location STAGING_LOCATION \ --temp-location TEMP_LOCATION \ --parameters PARAMETERS \ --additional-user-labels LABELS \
Sostituisci quanto segue:
JOB_ID: l'ID del jobTEMPLATE_FILE_LOCATION: il percorso di Cloud Storage del file di modelloREGION: la regione in cui eseguire il job DataflowSTAGING_LOCATION: il percorso di Cloud Storage per la gestione temporanea dei file localiTEMP_LOCATION: il percorso di Cloud Storage in cui scrivere i file temporanei. Se non viene impostato, il valore predefinito è il percorso di gestione temporanea.PARAMETERS: parametri della pipeline per il jobLABELS: (facoltativo) Etichette associate al job, utilizzando il formatoKEY_1=VALUE_1,KEY_2=VALUE_2,....
Durante il passaggio di gestione temporanea dell'avvio di un modello, Dataflow scrive i file nel percorso di gestione temporanea. Dataflow legge questi file di gestione temporanea per creare il grafico del job. Durante il passaggio di esecuzione, Dataflow scrive i file nella località temporanea.
Impostare le opzioni della pipeline
Per impostare le opzioni della pipeline quando esegui un modello flessibile, utilizza i seguenti flag nel comando gcloud dataflow flex-template run:
parameters: utilizza questo flag per impostare i seguenti tipi di opzioni della pipeline:Opzioni della pipeline supportate dai modelli flessibili. Per un elenco delle opzioni supportate dai modelli flessibili, consulta Opzioni della pipeline.
Opzioni della pipeline dichiarate in metadati del modello.
additional-pipeline-options: utilizza questo flag per impostare altre opzioni della pipeline Apache Beam che non sono supportate direttamente dai modelli flessibili.additional-experiments: Utilizza questo flag per impostare le opzioni della pipeline sperimentali (equivalenti all'opellaexperimentsopzione).
gcloud
Includi le opzioni della pipeline utilizzando il
parametersflag.Includi gli esperimenti di runtime e le opzioni della pipeline utilizzando i
additional-experimentse iadditional-pipeline-optionsflag.
Quando passi parametri o opzioni che contengono virgole o quando passi parametri di tipo List o Map, potresti dover eseguire una delle seguenti operazioni:
- Utilizza la
gcloud topic escapingsintassi per modificare i delimitatori. - Definisci i parametri in un file YAML e utilizza il
flags-fileflag. - Fornisci più flag
--additional-pipeline-optionso--additional-experimentscon l'opzione appropriata, ad esempio suggerimenti sulle risorse.
API
Includi le opzioni della pipeline utilizzando il
parameterscampo.Includi gli esperimenti di runtime e le opzioni della pipeline utilizzando i
additionalExperimentse iadditionalPipelineOptionscampi.
L'esempio seguente mostra come includere le opzioni della pipeline, gli esperimenti e le opzioni aggiuntive in un corpo della richiesta:
{
"jobName": "my-flex-template-job",
"parameters": {
"option_defined_in_metadata": "value"
},
"environment": {
"additionalExperiments": [
"use_runner_v2"
],
"additionalPipelineOptions": {
"common_pipeline_option": "value"
}
}
}
Quando utilizzi i modelli flessibili, puoi configurare alcune opzioni della pipeline durante l'inizializzazione della pipeline, ma altre opzioni della pipeline non possono essere modificate. Se gli argomenti della riga di comando richiesti dal modello flessibile vengono sovrascritti, il job potrebbe ignorare, sostituire o eliminare le opzioni della pipeline passate dal launcher del modello. Il job potrebbe non essere avviato o potrebbe essere avviato un job che non utilizza il modello flessibile. Per saperne di più, consulta Impossibile leggere il file del job.
Durante l'inizializzazione della pipeline, non modificare le seguenti opzioni della pipeline:
Java
runnerprojectjobNametemplateLocationregion
Python
runnerprojectjob_nametemplate_locationregion
Vai
runnerprojectjob_nametemplate_locationregion
Bloccare le chiavi SSH del progetto dalle VM che utilizzano chiavi SSH basate su metadati
Puoi impedire alle VM di accettare le chiavi SSH archiviate nei metadati del progetto bloccando le chiavi SSH del progetto dalle VM. Utilizza il flag additional-experiments con l'opzione di servizio block_project_ssh_keys:
--additional-experiments=block_project_ssh_keys
Per saperne di più, consulta Opzioni del servizio Dataflow.
Aggiornare un job del modello flessibile
La seguente richiesta di esempio mostra come aggiornare un job di streaming del modello utilizzando il projects.locations.flexTemplates.launch. Se vuoi utilizzare gcloud CLI, consulta Aggiornare una pipeline esistente.
Se vuoi aggiornare un modello classico, utilizza projects.locations.templates.launch.
Segui i passaggi per creare un job di streaming da un modello flessibile. Invia la seguente richiesta HTTP POST con i valori modificati:
POST https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/REGION/flexTemplates:launch { "launchParameter": { "update": true "jobName": "JOB_NAME", "parameters": { "input_subscription": "projects/PROJECT_ID/subscriptions/SUBSCRIPTION_NAME", "output_table": "PROJECT_ID:DATASET.TABLE_NAME" }, "containerSpecGcsPath": "STORAGE_PATH" }, }- Sostituisci
PROJECT_IDcon l'ID progetto. - Sostituisci
REGIONcon la regione Dataflow del job che stai aggiornando. - Sostituisci
JOB_NAMEcon il nome esatto del job che vuoi aggiornare. - Imposta
parameterssull'elenco di coppie chiave-valore. I parametri elencati sono specifici per questo esempio di modello. Se utilizzi un modello personalizzato, modifica i parametri in base alle esigenze. Se utilizzi il modello di esempio, sostituisci le seguenti variabili.- Sostituisci
SUBSCRIPTION_NAMEcon il nome dell'abbonamento Pub/Sub. - Sostituisci
DATASETcon il nome del set di dati BigQuery. - Sostituisci
TABLE_NAMEcon il nome della tabella BigQuery.
- Sostituisci
- Sostituisci
STORAGE_PATHcon il percorso di Cloud Storage del file di modello. Il percorso deve iniziare congs://.
- Sostituisci
Utilizza il parametro
environmentper modificare le impostazioni dell'ambiente. Per saperne di più, consultaFlexTemplateRuntimeEnvironment.(Facoltativo) Per inviare la richiesta utilizzando curl (Linux, macOS o Cloud Shell), salva la richiesta in un file JSON, quindi esegui il comando seguente:
curl -X POST -d "@FILE_PATH" -H "Content-Type: application/json" -H "Authorization: Bearer $(gcloud auth print-access-token)" https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/REGION/flexTemplates:launchSostituisci FILE_PATH con il percorso del file JSON contenente il corpo della richiesta.
Utilizza l' interfaccia di monitoraggio di Dataflow per verificare che sia stato creato un nuovo job con lo stesso nome. Questo job ha lo stato Aggiornato.
Passaggi successivi
- Scopri come creare un modello flessibile per la pipeline Apache Beam.
- Per saperne di più sui modelli classici, sui modelli flessibili e sui relativi scenari di utilizzo, consulta Modelli Dataflow.
- Per informazioni sulla risoluzione dei problemi relativi ai modelli flessibili, consulta Risolvere i problemi di timeout dei modelli flessibili.
- Per ulteriori architetture di riferimento, diagrammi e best practice, esplora il Cloud Architecture Center.