Il modello Generatore di dati di streaming genera record o messaggi sintetici e li invia a un sink di destinazione. Puoi configurare lo schema dei record e la velocità con cui vengono generati.
Il modello supporta i seguenti sink:
- Argomento Apache Kafka
- Tabella BigQuery
- Bucket Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Argomento Pub/Sub
- Tabella Spanner
Di seguito sono riportati alcuni possibili casi d'uso:
- Simula la pubblicazione di eventi in tempo reale su larga scala in un argomento Pub/Sub per misurare e determinare il numero e le dimensioni dei consumer necessari per elaborare gli eventi pubblicati.
- Genera dati sintetici per valutare i benchmark delle prestazioni o fungere da proof of concept.
- Convalida una pipeline end-to-end. Ad esempio, invia i record a un argomento Kafka, che vengono poi letti da un consumer downstream.
Requisiti della pipeline
- L'account di servizio worker
deve avere il ruolo Dataflow Worker (
roles/dataflow.worker) assegnato. Per maggiori informazioni, consulta Introduzione a IAM. - Il sink di destinazione deve esistere già.
- Configura il modello per il sink di destinazione. Per saperne di più, consulta Specifica il sink di destinazione.
- Definisci lo schema del record.
- Il modello genera dati JSON per impostazione predefinita. Per generare Avro o Parquet, vedi Specificare il formato di output.
Definisci lo schema del record
Il modello fornisce uno schema predefinito per i dati generati. Per utilizzare questo schema, imposta il parametro del modello schemaTemplate su GAME_EVENT.
In alternativa, puoi fornire il tuo schema dei dati nel seguente modo:
-
Crea un file di schema che contenga un modello JSON per i dati generati. Questo modello utilizza la libreria JSON Data Generator, che supporta varie funzioni per randomizzare i dati. Ad esempio:
{ "id": {{integer(0,1000)}}, "name": "{{uuid()}}", "isInStock": {{bool()}} }
Per saperne di più, consulta la documentazione di json-data-generator.
- Carica il file dello schema in un bucket Cloud Storage.
- Imposta il parametro del modello
schemaLocationsull'URI Cloud Storage del file modello.
Specifica il formato di output
Per impostazione predefinita, il modello produce dati JSON. Per alcune destinazioni, il modello supporta anche i formati Avro o Parquet:
- Avro: supportato per Cloud Storage, Apache Kafka e Pub/Sub
- Parquet: supportato per Cloud Storage.
Per generare l'output in formato Avro o Parquet:
- Imposta il parametro del modello
outputTypesuAVROper il formato Avro oPARQUETper il formato Parquet. - Crea un file di schema Avro.
- Carica il file dello schema su Cloud Storage.
- Imposta il parametro del modello
avroSchemaLocationsull'URI Cloud Storage del file di schema.
Specifica il sink di destinazione
Le sezioni seguenti descrivono come configurare il modello per ogni tipo di sink.
- Argomento Apache Kafka
- Tabella BigQuery
- Bucket Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Argomento Pub/Sub
- Tabella Spanner
Argomento Apache Kafka
Per scrivere in un argomento Kafka, imposta i seguenti parametri del modello:
sinkType:KAFKA.bootstrapServer: l'indirizzo bootstrap del cluster Kafka.kafkaTopic: L'argomento Kafka in cui scrivere.
Se stai scrivendo a un cluster Google Cloud Managed Service per Apache Kafka, concedi al
service account worker il ruolo Managed Kafka Client
(roles/managedkafka.client).
Tabella BigQuery
Per scrivere in una tabella BigQuery, imposta i seguenti parametri del modello:
sinkType:BIGQUERY.outputTableSpec: la tabella BigQuery in cui scrivere. Formatta questo parametro nel seguente modo:PROJECT_ID:DATASET.TABLE.
I seguenti parametri sono facoltativi:
outputDeadletterTable: il nome della tabella in cui la pipeline scrive i record non riusciti. Se non specificato, la pipeline crea una tabella denominataOUTPUT_TABLE_error_records, doveOUTPUT_TABLEè il nome della tabella di output.writeDisposition: specifica come scrivere in una tabella esistente. Sono supportati i seguenti valori:WRITE_APPEND. Aggiungi righe alla tabella esistente.WRITE_TRUNCATE. Tronca le righe esistenti.WRITE_EMPTY. Scrivi solo se la tabella è vuota. Se la tabella contiene già dati, il job non va a buon fine.
Il valore predefinito è
WRITE_APPEND.
Concedi al
service account worker il ruolo Editor dati BigQuery
(roles/bigquery.dataEditor).
Cloud Storage
Per scrivere in un bucket Cloud Storage, imposta i seguenti parametri del modello:
sinkType:GCS.outputDirectory: il percorso della cartella Cloud Storage in cui scrivere.
I seguenti parametri sono facoltativi:
numShards: il numero massimo di shard. Un valore più alto può consentire un throughput maggiore, ma potenzialmente costi di aggregazione dei dati più elevati. Se il valore è 0, Dataflow seleziona il numero di shard. Il valore predefinito è 0.outputFilenamePrefix: il prefisso del nome file. Il valore predefinito èoutput-.windowDuration: l'intervallo in cui la pipeline scrive i file in Cloud Storage. I formati consentiti sonoNs(secondi),Nm(minuti) eNh(ore). Il valore predefinito è1m(1 minuto).
Concedi al
service account worker il ruolo Amministratore oggetti Storage
(roles/storage.objectAdmin).
Endpoint JDBC
Per scrivere in un endpoint JDBC, imposta i seguenti parametri del modello:
sinkType:JDBC.driverClassName: la classe del driver JDBC da utilizzare. Esempio:com.mysql.jdbc.Driver.connectionUrl: la stringa di connessione per connettersi all'origine JDBC.-
statement: l'istruzione SQLINSERT INTOutilizzata per scrivere nel database. L'istruzione deve specificare le colonne della tabella da scrivere, con i caratteri segnaposto'?'per la clausolaVALUES. La pipeline sostituisce i segnaposto con i valori dei campi corrispondenti dei dati JSON.Esempio:
INSERT INTO tableName (column1, column2) VALUES (?,?).
I seguenti parametri sono facoltativi:
username: il nome utente per la connessione JDBC.password: la password per la connessione JDBC.connectionProperties: una stringa di proprietà per la connessione JDBC. Esempio:unicode=true;characterEncoding=UTF-8.
Argomento Pub/Sub
Per scrivere in un argomento Pub/Sub, imposta i seguenti parametri del modello:
sinkType:PUBSUB.topic: l'argomento Pub/Sub in cui scrivere.
Concedi al
service account worker il ruolo Publisher Pub/Sub
(roles/pubsub.publisher).
Tabella Spanner
Per scrivere in una tabella Spanner, imposta i seguenti parametri del modello:
sinkType:SPANNER.projectId: l'ID del progetto che contiene la tabella Spanner.spannerInstanceName: il nome dell'istanza Spanner.spannerDatabaseName: il nome del database Spanner.spannerTableName: il nome della tabella Spanner.
I seguenti parametri sono facoltativi:
maxNumMutations: il numero massimo di celle modificate per batch.maxNumRows: il numero massimo di righe modificate per batch.batchSizeBytes: il numero massimo di byte modificati per batch.commitDeadlineSeconds: il termine per la chiamata API commit, in secondi.
Concedi al
service account worker il ruolo Utente database Cloud Spanner
(roles/spanner.databaseUser).
Parametri del modello
Parametri obbligatori
- qps: indica la velocità di pubblicazione dei messaggi al secondo su Pub/Sub.
Parametri facoltativi
- schemaTemplate: modello di schema preesistente da utilizzare. Il valore deve essere uno dei seguenti: [GAME_EVENT].
- schemaLocation: percorso Cloud Storage della posizione dello schema. Ad esempio,
gs://<bucket-name>/prefix. - topic: il nome dell'argomento a cui la pipeline deve pubblicare i dati. Ad esempio,
projects/<project-id>/topics/<topic-name>. - messagesLimit: indica il numero massimo di messaggi di output da generare. 0 significa illimitato. Il valore predefinito è 0.
- outputType: il tipo di output del messaggio. Il valore predefinito è JSON.
- avroSchemaLocation: percorso Cloud Storage della posizione dello schema Avro. Obbligatorio quando il tipo di output è AVRO o PARQUET. Ad esempio,
gs://your-bucket/your-path/schema.avsc. - sinkType: il tipo di sink del messaggio. Il valore predefinito è PUBSUB.
- outputTableSpec: tabella BigQuery di output. Obbligatorio quando sinkType è BIGQUERY. Ad esempio,
<project>:<dataset>.<table_name>. - writeDisposition: BigQuery WriteDisposition. Ad esempio, WRITE_APPEND, WRITE_EMPTY o WRITE_TRUNCATE. Il valore predefinito è WRITE_APPEND.
- outputDeadletterTable: i messaggi che non sono riusciti a raggiungere la tabella di output per vari motivi (ad es. schema non corrispondente, JSON non valido) vengono scritti in questa tabella. Se non esiste, verrà creato durante l'esecuzione della pipeline. Ad esempio,
your-project-id:your-dataset.your-table-name. - windowDuration: la durata/dimensione della finestra in cui i dati verranno scritti in Cloud Storage. I formati consentiti sono: Ns (per i secondi, ad esempio 5s), Nm (per i minuti, ad esempio 12m), Nh (per le ore, ad esempio 2h). Ad esempio,
1m. Il valore predefinito è 1 minuto. - outputDirectory: il percorso e il prefisso del nome file per la scrittura dei file di output. Deve terminare con una barra. La formattazione DateTime viene utilizzata per analizzare il percorso della directory per i formattatori di data e ora. Ad esempio,
gs://your-bucket/your-path/. - outputFilenamePrefix: il prefisso da inserire in ogni file in finestra. Ad esempio,
output-. Il valore predefinito è: output-. - numShards: il numero massimo di shard di output prodotti durante la scrittura. Un numero maggiore di shard significa una velocità effettiva più elevata per la scrittura in Cloud Storage, ma potenzialmente un costo di aggregazione dei dati più elevato tra gli shard durante l'elaborazione dei file Cloud Storage di output. Il valore predefinito viene deciso da Dataflow.
- driverClassName: nome della classe del driver JDBC da utilizzare. Ad esempio,
com.mysql.jdbc.Driver. - connectionUrl: stringa di connessione dell'URL per la connessione all'origine JDBC. Ad esempio,
jdbc:mysql://some-host:3306/sampledb. - username: nome utente da utilizzare per la connessione JDBC.
- password: password da utilizzare per la connessione JDBC.
- connectionProperties: stringa di proprietà da utilizzare per la connessione JDBC. Il formato della stringa deve essere [propertyName=property;]*. Ad esempio,
unicode=true;characterEncoding=UTF-8. - statement: istruzione SQL che verrà eseguita per scrivere nel database. L'istruzione deve specificare i nomi delle colonne della tabella in qualsiasi ordine. Dal file JSON verranno letti e aggiunti all'istruzione solo i valori dei nomi delle colonne specificati. Ad esempio,
INSERT INTO tableName (column1, column2) VALUES (?,?). - projectId: l'ID progetto Google Cloud in cui si trova la tabella Spanner.
- spannerInstanceName: nome dell'istanza Cloud Spanner.
- spannerDatabaseName: nome del database Cloud Spanner.
- spannerTableName: nome della tabella Cloud Spanner.
- maxNumMutations: specifica il limite di mutazione delle celle (numero massimo di celle mutate per batch). Il valore predefinito è 5000.
- maxNumRows: specifica il limite di mutazione delle righe (numero massimo di righe mutate per batch). Il valore predefinito è 1000.
- batchSizeBytes: specifica il limite di dimensione del batch (numero massimo di byte modificati per batch). Il valore predefinito è 1 MB.
- commitDeadlineSeconds: specifica la scadenza in secondi per la chiamata API Commit.
- bootstrapServer: server bootstrap Kafka. Ad esempio,
localhost:9092. - kafkaTopic: argomento Kafka in cui scrivere. Ad esempio,
topic.
Esegui il modello
Console
- Vai alla pagina Crea job da modello di Dataflow. Vai a Crea job da modello
- Nel campo Nome job, inserisci un nome univoco per il job.
- (Facoltativo) Per Endpoint a livello di regione, seleziona un valore dal menu a discesa. La regione
predefinita è
us-central1.Per un elenco delle regioni in cui puoi eseguire un job Dataflow, consulta Località di Dataflow.
- Dal menu a discesa Modello di dataflow, seleziona il modello Generatore di dati di streaming.
- Nei campi dei parametri forniti, inserisci i valori dei parametri.
- Fai clic su Esegui job.
gcloud
Nella shell o nel terminale, esegui il modello:
gcloud dataflow flex-template run JOB_NAME \ --project=PROJECT_ID \ --region=REGION_NAME \ --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/ \ --parameters \ schemaLocation=SCHEMA_LOCATION,\ qps=QPS,\ topic=PUBSUB_TOPIC
Sostituisci quanto segue:
PROJECT_ID: l'ID progetto Google Cloud in cui vuoi eseguire il job DataflowREGION_NAME: la regione in cui vuoi eseguire il deployment del job Dataflow, ad esempious-central1JOB_NAME: un nome univoco del job a tua sceltaVERSION: la versione del modello che vuoi utilizzarePuoi utilizzare i seguenti valori:
latestper utilizzare l'ultima versione del modello, disponibile nella cartella principale senza data nel bucket: gs://dataflow-templates-REGION_NAME/latest/- il nome della versione, ad esempio
2023-09-12-00_RC00, per utilizzare una versione specifica del modello, che si trova nidificata nella rispettiva cartella principale con data nel bucket: gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: il percorso del file schema in Cloud Storage. Ad esempio:gs://mybucket/filename.json.QPS: il numero di messaggi da pubblicare al secondoPUBSUB_TOPIC: l'argomento Pub/Sub di output. Ad esempio:projects/my-project-id/topics/my-topic-id.
API
Per eseguire il modello utilizzando l'API REST, invia una richiesta POST HTTP. Per saperne di più sull'API e sui relativi ambiti di autorizzazione, consulta projects.templates.launch.
POST https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/LOCATION/flexTemplates:launch { "launch_parameter": { "jobName": "JOB_NAME", "parameters": { "schemaLocation": "SCHEMA_LOCATION", "qps": "QPS", "topic": "PUBSUB_TOPIC" }, "containerSpecGcsPath": "gs://dataflow-templates-LOCATION/VERSION/flex/", } }
Sostituisci quanto segue:
PROJECT_ID: l'ID progetto Google Cloud in cui vuoi eseguire il job DataflowLOCATION: la regione in cui vuoi eseguire il deployment del job Dataflow, ad esempious-central1JOB_NAME: un nome univoco del job a tua sceltaVERSION: la versione del modello che vuoi utilizzarePuoi utilizzare i seguenti valori:
latestper utilizzare l'ultima versione del modello, disponibile nella cartella principale senza data nel bucket: gs://dataflow-templates-REGION_NAME/latest/- il nome della versione, ad esempio
2023-09-12-00_RC00, per utilizzare una versione specifica del modello, che si trova nidificata nella rispettiva cartella principale con data nel bucket: gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: il percorso del file schema in Cloud Storage. Ad esempio:gs://mybucket/filename.json.QPS: il numero di messaggi da pubblicare al secondoPUBSUB_TOPIC: l'argomento Pub/Sub di output. Ad esempio:projects/my-project-id/topics/my-topic-id.
Passaggi successivi
- Scopri di più sui modelli Dataflow.
- Consulta l'elenco dei modelli forniti da Google.