La plantilla de Generador de datos de transmisión genera registros o mensajes sintéticos y los envía a un receptor de destino. Puedes configurar el esquema de registros y la frecuencia con la que se generan.
La plantilla admite los siguientes receptores:
- Tema de Apache Kafka
- Tabla de BigQuery
- Bucket de Cloud Storage
- Endpoint de conectividad a bases de datos de Java (JDBC)
- Tema de Pub/Sub
- Tabla de Spanner
A continuación, se incluyen algunos casos de uso posibles:
- Simula la publicación de eventos en tiempo real a gran escala en un tema de Pub/Sub para medir y determinar la cantidad y el tamaño de los consumidores necesarios para procesar los eventos publicados.
- Genera datos sintéticos para evaluar comparativas de rendimiento o servir como prueba de concepto.
- Valida una canalización de extremo a extremo. Por ejemplo, enviar registros a un tema de Kafka, que luego lee un consumidor de nivel inferior
Requisitos de la canalización
- La cuenta de servicio de trabajador necesita la función asignada de trabajador de Dataflow (
roles/dataflow.worker). Para obtener más información, consulta Introducción a IAM. - El receptor de destino ya debe existir.
- Configura la plantilla para el receptor de destino. Para obtener más información, consulta Cómo especificar el receptor de destino.
- Define el esquema del registro.
- De forma predeterminada, la plantilla genera datos JSON. Para generar Avro o Parquet, consulta Cómo especificar el formato de salida.
Define el esquema del registro
La plantilla proporciona un esquema predefinido para los datos generados. Para usar este esquema, establece el parámetro de plantilla schemaTemplate en GAME_EVENT.
Como alternativa, puedes proporcionar tu propio esquema de datos de la siguiente manera:
-
Crea un archivo de esquema que contenga una plantilla JSON para los datos generados. Esta plantilla usa la biblioteca del Generador de datos JSON, que admite varias funciones para aleatorizar los datos. Por ejemplo:
{ "id": {{integer(0,1000)}}, "name": "{{uuid()}}", "isInStock": {{bool()}} }
Para obtener más información, consulta la documentación de json-data-generator.
- Sube el archivo de esquema a un bucket de Cloud Storage.
- Establece el parámetro de plantilla
schemaLocationen el URI de Cloud Storage del archivo de plantilla.
Especifica el formato del resultado
De forma predeterminada, la plantilla genera datos JSON. En el caso de algunos destinos, la plantilla también admite los formatos Avro o Parquet:
- Avro: Se admite en Cloud Storage, Apache Kafka y Pub/Sub
- Parquet: Se admite en Cloud Storage.
Para generar el formato Avro o Parquet, haz lo siguiente:
- Establece el parámetro de plantilla
outputTypeenAVROpara el formato Avro o enPARQUETpara el formato Parquet. - Crea un archivo de esquema de Avro.
- Sube el archivo de esquema a Cloud Storage.
- Establece el parámetro de plantilla
avroSchemaLocationen el URI de Cloud Storage del archivo de esquema.
Especifica el receptor de destino
En las siguientes secciones, se describe cómo configurar la plantilla para cada tipo de receptor.
- Tema de Apache Kafka
- Tabla de BigQuery
- Bucket de Cloud Storage
- Endpoint de Java Database Connectivity (JDBC)
- Tema de Pub/Sub
- Tabla de Spanner
Tema de Apache Kafka
Para escribir en un tema de Kafka, establece los siguientes parámetros de plantilla:
sinkType:KAFKA.bootstrapServer: Es la dirección de arranque del clúster de Kafka.kafkaTopic: Es el tema de Kafka en el que se escribirá.
Si escribes en un clúster de Google Cloud Managed Service para Apache Kafka, otorga a la
cuenta de servicio del trabajador el rol de cliente de Kafka administrado (roles/managedkafka.client).
Tabla de BigQuery
Para escribir en una tabla de BigQuery, configura los siguientes parámetros de plantilla:
sinkType:BIGQUERY.outputTableSpec: Es la tabla de BigQuery en la que se escribirán los datos. Da formato a este parámetro de la siguiente manera:PROJECT_ID:DATASET.TABLE.
Los parámetros que aparecen a continuación son opcionales:
outputDeadletterTable: Es el nombre de la tabla en la que la canalización escribe los registros con errores. Si no se especifica, la canalización crea una tabla llamadaOUTPUT_TABLE_error_records, en la queOUTPUT_TABLEes el nombre de la tabla de salida.writeDisposition: Especifica cómo escribir en una tabla existente. Se admiten los siguientes valores:WRITE_APPEND: Agrega filas a la tabla existente.WRITE_TRUNCATE. Trunca las filas existentes.WRITE_EMPTY. Escribe solo si la tabla está vacía. Si la tabla ya tiene datos, el trabajo falla.
El valor predeterminado es
WRITE_APPEND.
Otorga a la
cuenta de servicio del trabajador el rol de editor de datos de BigQuery
(roles/bigquery.dataEditor).
Cloud Storage
Para escribir en un bucket de Cloud Storage, establece los siguientes parámetros de plantilla:
sinkType:GCS.outputDirectory: Es la ruta de acceso a la carpeta de Cloud Storage en la que se escribirán los datos.
Los parámetros que aparecen a continuación son opcionales:
numShards: Es la cantidad máxima de fragmentos. Un valor más alto puede permitir una mayor capacidad de procesamiento, pero también puede generar costos más altos de agregación de datos. Si el valor es 0, Dataflow selecciona la cantidad de fragmentos. El valor predeterminado es 0.outputFilenamePrefix: Es el prefijo del nombre de archivo. El valor predeterminado esoutput-.windowDuration: Es el intervalo en el que la canalización escribe archivos en Cloud Storage. Los formatos permitidos sonNs(segundos),Nm(minutos) yNh(horas). El valor predeterminado es1m(1 minuto).
Otorga a la
cuenta de servicio del trabajador el rol de administrador de objetos de almacenamiento (roles/storage.objectAdmin).
Extremo de JDBC
Para escribir en un extremo de JDBC, establece los siguientes parámetros de plantilla:
sinkType:JDBC.driverClassName: Es la clase del controlador de JDBC que se usará. Ejemplo:com.mysql.jdbc.Driver.connectionUrl: Es la cadena de conexión para conectarse a la fuente de JDBC.-
statement: Es la instrucción de SQLINSERT INTOque se usó para escribir en la base de datos. La instrucción debe especificar qué columnas de la tabla se escribirán, con caracteres de marcador de posición'?'para la cláusulaVALUES. La canalización reemplaza los marcadores de posición por los valores de campo correspondientes de los datos JSON.Ejemplo:
INSERT INTO tableName (column1, column2) VALUES (?,?).
Los parámetros que aparecen a continuación son opcionales:
username: Es el nombre de usuario para la conexión de JDBC.password: Es la contraseña de la conexión de JDBC.connectionProperties: Es una cadena de propiedades para la conexión de JDBC. Ejemplo:unicode=true;characterEncoding=UTF-8.
Tema de Pub/Sub
Para escribir en un tema de Pub/Sub, establece los siguientes parámetros de plantilla:
sinkType:PUBSUB.topic: Es el tema de Pub/Sub en el que se escribirá.
Otorga a la
cuenta de servicio del trabajador el rol de publicador de Pub/Sub
(roles/pubsub.publisher).
Tabla de Spanner
Para escribir en una tabla de Spanner, establece los siguientes parámetros de plantilla:
sinkType:SPANNER.projectId: Es el ID del proyecto que contiene la tabla de Spanner.spannerInstanceName: Es el nombre de la instancia de Spanner.spannerDatabaseName: Es el nombre de la base de datos de Spanner.spannerTableName: Es el nombre de la tabla de Spanner.
Los parámetros que aparecen a continuación son opcionales:
maxNumMutations: Es la cantidad máxima de celdas modificadas por lote.maxNumRows: Es la cantidad máxima de filas modificadas por lote.batchSizeBytes: Es la cantidad máxima de bytes que se modifican por lote.commitDeadlineSeconds: Es la fecha límite para la llamada a la API de confirmación, en segundos.
Otorga a la
cuenta de servicio del trabajador el rol de Usuario de bases de datos de Cloud Spanner (roles/spanner.databaseUser).
Parámetros de la plantilla
Parámetros obligatorios
- qps: Indica la tasa de mensajes por segundo que se publicarán en Pub/Sub.
Parámetros opcionales
- schemaTemplate: Es la plantilla de esquema preexistente que se usará. El valor debe ser uno de los siguientes: [GAME_EVENT].
- schemaLocation: Es la ruta de acceso de Cloud Storage a la ubicación del esquema. Por ejemplo,
gs://<bucket-name>/prefix - topic: Es el nombre del tema en el que la canalización debe publicar datos. Por ejemplo,
projects/<project-id>/topics/<topic-name> - messagesLimit: Indica la cantidad máxima de mensajes de salida que se generarán. El valor 0 significa ilimitado. La configuración predeterminada es 0.
- outputType: Es el tipo de resultado del mensaje. El valor predeterminado es JSON.
- avroSchemaLocation: Es la ruta de Cloud Storage a la ubicación del esquema de Avro. Es obligatorio cuando el tipo de salida es AVRO o PARQUET. Por ejemplo,
gs://your-bucket/your-path/schema.avsc - sinkType: Es el tipo de receptor de mensajes. El valor predeterminado es PUBSUB.
- outputTableSpec: Es la tabla de salida de BigQuery. Obligatorio cuando sinkType es BIGQUERY. Por ejemplo,
<project>:<dataset>.<table_name>. - writeDisposition: WriteDisposition de BigQuery. Por ejemplo, WRITE_APPEND, WRITE_EMPTY o WRITE_TRUNCATE. La configuración predeterminada es: WRITE_APPEND.
- outputDeadletterTable: Los mensajes que no llegaron a la tabla de resultados por diversos motivos (p.ej., un esquema no coincidente, un archivo JSON con formato incorrecto) se escriben en esta tabla. Si no existe, se creará durante la ejecución de la canalización. Por ejemplo,
your-project-id:your-dataset.your-table-name - windowDuration: Es la duración o el tamaño de la ventana en la que se escribirán los datos en Cloud Storage. Los formatos permitidos son: Ns (para los segundos, por ejemplo, 5 s), Nm (para los minutos, por ejemplo, 12 m) y Nh (para las horas, por ejemplo, 2 h). Por ejemplo,
1mLa configuración predeterminada es de 1 m. - outputDirectory: Es la ruta de acceso y el prefijo del nombre de archivo para escribir los archivos de salida. Debe terminar con una barra. El formato de fecha y hora se usa a fin de analizar la ruta del directorio para los formateadores de fecha y hora. Por ejemplo,
gs://your-bucket/your-path/ - outputFilenamePrefix: Es el prefijo para colocar en cada archivo con ventanas. Por ejemplo,
output-La configuración predeterminada es: output-. - numShards: Es la cantidad máxima de fragmentos de salida que se produce con la escritura. Una mayor cantidad de fragmentos implica una mayor capacidad de procesamiento para la escritura en Cloud Storage, pero, también, un mayor costo de agregación de datos entre fragmentos cuando se procesan archivos de salida de Cloud Storage. Dataflow decide el valor predeterminado.
- driverClassName: Es el nombre de la clase del controlador JDBC que se usará. Por ejemplo,
com.mysql.jdbc.Driver - connectionUrl: Cadena de conexión de URL para conectarse a la fuente de JDBC. Por ejemplo,
jdbc:mysql://some-host:3306/sampledb - username: Nombre de usuario que se usará para la conexión de JDBC.
- password: Contraseña para usar en la conexión de JDBC.
- connectionProperties: Cadena de propiedades para usar en la conexión de JDBC. El formato de la cadena debe ser [propertyName=property;]*. Por ejemplo,
unicode=true;characterEncoding=UTF-8 - statement: Es la instrucción de SQL que se ejecutará para escribir en la base de datos. La instrucción debe especificar los nombres de columna de la tabla en cualquier orden. Solo se leerán los valores de los nombres de columna especificados desde el archivo JSON y se agregarán a la instrucción. Por ejemplo,
INSERT INTO tableName (column1, column2) VALUES (?,?) - projectId: Es el ID del proyecto de GCP en el que se encuentra la tabla de Spanner.
- spannerInstanceName: Es el nombre de la instancia de Cloud Spanner.
- spannerDatabaseName: Es el nombre de la base de datos de Cloud Spanner.
- spannerTableName: Es el nombre de la tabla de Cloud Spanner.
- maxNumMutations: Especifica el límite de mutación de celdas (cantidad máxima de celdas mutadas por lote). El valor predeterminado es 5,000.
- maxNumRows: Especifica el límite de mutación de filas (cantidad máxima de filas mutadas por lote). El valor predeterminado es 1,000.
- batchSizeBytes: Especifica el límite de tamaño del lote (cantidad máxima de bytes mutados por lote). El valor predeterminado es 1 MB.
- commitDeadlineSeconds: Especifica el plazo en segundos para la llamada a la API de Commit.
- bootstrapServer: Servidor de arranque de Kafka. Por ejemplo,
localhost:9092. - kafkaTopic: Tema de Kafka en el que se escribirá. Por ejemplo,
topic
Ejecuta la plantilla
Console
- Ve a la página Crear un trabajo a partir de una plantilla de Dataflow. Ir a Crear un trabajo a partir de una plantilla
- En el campo Nombre del trabajo, ingresa un nombre de trabajo único.
- Opcional: Para Extremo regional, selecciona un valor del menú desplegable. La región predeterminada es
us-central1.Para obtener una lista de regiones en las que puedes ejecutar un trabajo de Dataflow, consulta Ubicaciones de Dataflow.
- En el menú desplegable Plantilla de Dataflow, selecciona la plantilla Streaming Data Generator.
- En los campos de parámetros proporcionados, ingresa los valores de tus parámetros.
- Haga clic en Ejecutar trabajo.
gcloud
En tu shell o terminal, ejecuta la plantilla:
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
Reemplaza lo siguiente:
PROJECT_ID: El ID del proyecto Google Cloud en el que deseas ejecutar el trabajo de DataflowREGION_NAME: La región en la que deseas implementar tu trabajo de Dataflow, por ejemplo,us-central1JOB_NAME: Es el nombre del trabajo que elijasVERSION: Es la versión de la plantilla que deseas usar.Puedes usar los siguientes valores:
latestpara usar la última versión de la plantilla, que está disponible en la carpeta superior non-dated en el bucket gs://dataflow-templates-REGION_NAME/latest/- el nombre de la versión, como
2023-09-12-00_RC00, para usar una versión específica de la plantilla, que se puede encontrar anidada en la carpeta superior con fecha correspondiente en el bucket gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: Es la ruta de acceso al archivo de esquema en Cloud Storage. Por ejemplo:gs://mybucket/filename.json.QPS: Es la cantidad de mensajes que se publicarán por segundo.PUBSUB_TOPIC: Es el tema de salida de Pub/Sub. Por ejemplo:projects/my-project-id/topics/my-topic-id.
API
Para ejecutar la plantilla con la API de REST, envía una solicitud POST HTTP. Para obtener más información de la API y sus permisos de autorización, 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/", } }
Reemplaza lo siguiente:
PROJECT_ID: El ID del proyecto Google Cloud en el que deseas ejecutar el trabajo de DataflowLOCATION: La región en la que deseas implementar tu trabajo de Dataflow, por ejemplo,us-central1JOB_NAME: Es el nombre del trabajo que elijasVERSION: Es la versión de la plantilla que deseas usar.Puedes usar los siguientes valores:
latestpara usar la última versión de la plantilla, que está disponible en la carpeta superior non-dated en el bucket gs://dataflow-templates-REGION_NAME/latest/- el nombre de la versión, como
2023-09-12-00_RC00, para usar una versión específica de la plantilla, que se puede encontrar anidada en la carpeta superior con fecha correspondiente en el bucket gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: Es la ruta de acceso al archivo de esquema en Cloud Storage. Por ejemplo:gs://mybucket/filename.json.QPS: Es la cantidad de mensajes que se publicarán por segundo.PUBSUB_TOPIC: Es el tema de salida de Pub/Sub. Por ejemplo:projects/my-project-id/topics/my-topic-id.
¿Qué sigue?
- Obtén información sobre las plantillas de Dataflow.
- Consulta la lista de plantillas que proporciona Google.