O modelo Gerador de dados de streaming gera registros sintéticos ou mensagens e os envia para um coletor de destino. É possível configurar o esquema de registro e a taxa de geração de registros.
O modelo é compatível com os seguintes coletores:
- Tópico do Apache Kafka
- Tabela do BigQuery
- Bucket do Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Tópico do Pub/Sub
- Tabela do Spanner
Veja a seguir um conjunto de alguns casos de uso possíveis:
- Simule a publicação de eventos em tempo real em grande escala em um tópico do Pub/Sub para medir e determinar o número e o tamanho dos consumidores necessários para processar eventos publicados.
- Gere dados sintéticos para avaliar comparativos de mercado de desempenho ou servir como prova de conceito.
- Validar um pipeline de ponta a ponta. Por exemplo, envie registros para um tópico do Kafka, que são lidos por um consumidor downstream.
Requisitos de pipeline
- A conta de serviço do worker precisa do papel atribuído do Dataflow Worker (
roles/dataflow.worker). Para mais informações, consulte Introdução ao IAM. - O coletor de destino já precisa existir.
- Configure o modelo para o coletor de destino. Para mais informações, consulte Especificar o gravador de destino.
- Defina o esquema de registro.
- Por padrão, o modelo gera dados JSON. Para gerar Avro ou Parquet, consulte Especificar o formato de saída.
Definir o esquema de registro
O modelo fornece um esquema predefinido para os dados gerados. Para usar esse
esquema, defina o parâmetro de modelo schemaTemplate como
GAME_EVENT.
Como alternativa, você pode fornecer seu próprio esquema de dados da seguinte forma:
-
Crie um arquivo de esquema que contenha um modelo JSON para os dados gerados. Esse modelo usa a biblioteca Gerador de dados JSON, que é compatível com várias funções para aleatorizar os dados. Exemplo:
{ "id": {{integer(0,1000)}}, "name": "{{uuid()}}", "isInStock": {{bool()}} }
Para mais informações, consulte a documentação do json-data-generator.
- Faça upload do arquivo de esquema para um bucket do Cloud Storage.
- Defina o parâmetro de modelo
schemaLocationcomo o URI do Cloud Storage do arquivo de modelo.
especifique o formato das respostas
Por padrão, o modelo gera dados JSON. Para alguns destinos, o modelo também é compatível com os formatos Avro ou Parquet:
- Avro: compatível com Cloud Storage, Apache Kafka e Pub/Sub
- Parquet: compatível com o Cloud Storage.
Para gerar no formato Avro ou Parquet, faça o seguinte:
- Defina o parâmetro de modelo
outputTypecomoAVROpara o formato Avro ouPARQUETpara o formato Parquet. - Crie um arquivo de esquema Avro.
- Faça upload do arquivo de esquema para o Cloud Storage.
- Defina o parâmetro de modelo
avroSchemaLocationcomo o URI do Cloud Storage do arquivo de esquema.
Especificar o coletor de destino
As seções a seguir descrevem como configurar o modelo para cada tipo de coletor.
- Tópico do Apache Kafka
- Tabela do BigQuery
- Bucket do Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Tópico do Pub/Sub
- Tabela do Spanner
Tópico do Apache Kafka
Para gravar em um tópico do Kafka, defina os seguintes parâmetros de modelo:
sinkType:KAFKA.bootstrapServer: o endereço de inicialização do cluster do Kafka.kafkaTopic: o tópico do Kafka em que gravar.
Se você estiver gravando em um cluster do Serviço Gerenciado do Google Cloud para Apache Kafka, conceda à
conta de serviço do worker a função de cliente do Kafka gerenciado (roles/managedkafka.client).
Tabela do BigQuery
Para gravar em uma tabela do BigQuery, defina os seguintes parâmetros de modelo:
sinkType:BIGQUERY.outputTableSpec: a tabela do BigQuery em que será gravada. Formate esse parâmetro da seguinte maneira:PROJECT_ID:DATASET.TABLE.
Os seguintes parâmetros são opcionais:
outputDeadletterTable: o nome da tabela em que o pipeline grava registros com falha. Se não for especificado, o pipeline vai criar uma tabela chamadaOUTPUT_TABLE_error_records, em queOUTPUT_TABLEé o nome da tabela de saída.writeDisposition: especifica como gravar em uma tabela existente. Os valores a seguir são compatíveis:WRITE_APPEND. Anexe linhas à tabela atual.WRITE_TRUNCATE. Truncar as linhas atuais.WRITE_EMPTY. Gravar apenas se a tabela estiver vazia. Se a tabela já tiver dados, o job vai falhar.
O valor padrão é
WRITE_APPEND.
Conceda à
conta de serviço do worker o papel de editor de dados do BigQuery (roles/bigquery.dataEditor).
Cloud Storage
Para gravar em um bucket do Cloud Storage, defina os seguintes parâmetros de modelo:
sinkType:GCS.outputDirectory: o caminho da pasta do Cloud Storage para gravar.
Os seguintes parâmetros são opcionais:
numShards: o número máximo de fragmentos. Um valor maior pode permitir uma capacidade de processamento maior, mas custos de agregação de dados potencialmente mais altos. Se o valor for 0, o Dataflow vai selecionar o número de fragmentos. O valor padrão é 0.outputFilenamePrefix: o prefixo do nome do arquivo. O valor padrão éoutput-.windowDuration: o intervalo em que o pipeline grava arquivos no Cloud Storage. Os formatos permitidos sãoNs(segundos),Nm(minutos) eNh(horas). O valor padrão é1m(1 minuto).
Conceda à
conta de serviço do worker o papel Administrador de objetos do Storage (roles/storage.objectAdmin).
Endpoint JDBC
Para gravar em um endpoint JDBC, defina os seguintes parâmetros de modelo:
sinkType:JDBC.driverClassName: a classe do driver JDBC a ser usada. Exemplo:com.mysql.jdbc.Driver.connectionUrl: a string de conexão para acessar a fonte JDBC.-
statement: a instrução SQLINSERT INTOusada para gravar no banco de dados. A instrução precisa especificar quais colunas da tabela serão gravadas, com caracteres de marcador de posição'?'para a cláusulaVALUES. O pipeline substitui os marcadores pelos valores de campo correspondentes dos dados JSON.Exemplo:
INSERT INTO tableName (column1, column2) VALUES (?,?).
Os seguintes parâmetros são opcionais:
username: o nome de usuário da conexão JDBC.password: a senha da conexão JDBC.connectionProperties: uma string de propriedades para a conexão JDBC. Exemplo:unicode=true;characterEncoding=UTF-8.
Tópico do Pub/Sub
Para gravar em um tópico do Pub/Sub, defina os seguintes parâmetros de modelo:
sinkType:PUBSUB.topic: o tópico do Pub/Sub em que gravar.
Conceda à
conta de serviço do worker o papel de Publicador do Pub/Sub
(roles/pubsub.publisher).
Tabela do Spanner
Para gravar em uma tabela do Spanner, defina os seguintes parâmetros do modelo:
sinkType:SPANNER.projectId: o ID do projeto que contém a tabela do Spanner.spannerInstanceName: o nome da instância do Spanner.spannerDatabaseName: o nome do banco de dados do Spanner.spannerTableName: o nome da tabela do Spanner.
Os seguintes parâmetros são opcionais:
maxNumMutations: o número máximo de células mutantes por lote.maxNumRows: o número máximo de linhas mutadas por lote.batchSizeBytes: o número máximo de bytes mutados por lote.commitDeadlineSeconds: o prazo para a chamada de API de commit, em segundos.
Conceda à
conta de serviço do worker o papel de Usuário do banco de dados do Cloud Spanner (roles/spanner.databaseUser).
Parâmetros do modelo
Parâmetros obrigatórios
- qps: indica a taxa de mensagens por segundo a serem publicadas no Pub/Sub.
Parâmetros opcionais
- schemaTemplate: modelo de esquema preexistente a ser usado. O valor precisa ser um destes: [GAME_EVENT].
- schemaLocation: caminho do Cloud Storage para o local do esquema. Por exemplo,
gs://<bucket-name>/prefix. - topic: o nome do tópico em que o pipeline precisa publicar dados. Por exemplo,
projects/<project-id>/topics/<topic-name>. - messagesLimit: indica o número máximo de mensagens de saída a serem geradas. 0 significa ilimitado. Padrão: 0.
- outputType: o tipo de saída da mensagem. O padrão é JSON.
- avroSchemaLocation: caminho do Cloud Storage para o local do esquema Avro. Obrigatório quando o tipo de saída é AVRO ou PARQUET. Por exemplo,
gs://your-bucket/your-path/schema.avsc. - sinkType: o tipo de coletor de mensagens. O padrão é PUBSUB.
- outputTableSpec: tabela do BigQuery de saída. Obrigatório quando "sinkType" é BIGQUERY. Por exemplo,
<project>:<dataset>.<table_name>. - writeDisposition: WriteDisposition do BigQuery. Por exemplo, WRITE_APPEND, WRITE_EMPTY ou WRITE_TRUNCATE. O valor padrão é: WRITE_APPEND.
- outputDeadletterTable: mensagens que não alcançaram a tabela de saída por vários motivos (por exemplo, esquema incompatível, json incorreto) são gravadas nesta tabela. Se não existir, será criada durante a execução do pipeline. Por exemplo,
your-project-id:your-dataset.your-table-name. - windowDuration: a duração/tamanho da janela em que os dados serão gravados no Cloud Storage. Os formatos permitidos são: Ns (para segundos, exemplo: "5s"), Nm (para minutos, exemplo: "12m") e Nh (para horas, exemplo: "2h"). Por exemplo,
1m. O padrão é de 1 minuto. - outputDirectory: o caminho e o prefixo do nome do arquivo para gravar arquivos de saída. Precisa terminar com uma barra. A formatação DateTime é usada para analisar o caminho do diretório em busca de formatadores de data e hora. Por exemplo,
gs://your-bucket/your-path/. - outputFilenamePrefix: o prefixo a ser colocado em cada arquivo em janela. Por exemplo,
output-. O padrão é: output-. - numShards: o número máximo de fragmentos de saída produzidos durante a gravação. Um número maior de fragmentos significa maior capacidade de gravação no Cloud Storage, mas um custo de agregação de dados potencialmente maior entre os fragmentos ao processar os arquivos de saída do Cloud Storage. O valor padrão é decidido pelo Dataflow.
- driverClassName: nome da classe do driver JDBC a ser usado. Por exemplo,
com.mysql.jdbc.Driver. - connectionUrl: string de conexão de URL para acessar a fonte JDBC. Por exemplo,
jdbc:mysql://some-host:3306/sampledb. - username: nome de usuário a ser usado para a conexão JDBC.
- password: senha a ser usada para a conexão JDBC.
- connectionProperties: string de propriedades a ser usada para a conexão JDBC. O formato da string precisa ser [propertyName=property;]*. Por exemplo,
unicode=true;characterEncoding=UTF-8. - statement: instrução SQL que será executada para gravar no banco de dados. A instrução precisa especificar os nomes das colunas da tabela em qualquer ordem. Somente os valores dos nomes das colunas especificadas serão lidos no JSON e adicionados à instrução. Por exemplo,
INSERT INTO tableName (column1, column2) VALUES (?,?). - projectId: ID do projeto do GCP em que a tabela do Spanner está localizada.
- spannerInstanceName: nome da instância do Cloud Spanner.
- spannerDatabaseName: nome do banco de dados do Cloud Spanner.
- spannerTableName: nome da tabela do Cloud Spanner.
- maxNumMutations: especifica o limite de mutação de células (número máximo de células mutadas por lote). O valor padrão é 5.000.
- maxNumRows: especifica o limite de mutação de linhas (número máximo de linhas mutadas por lote). O valor padrão é 1000.
- batchSizeBytes: especifica o limite de tamanho do lote (número máximo de bytes mutados por lote). O valor padrão é 1 MB.
- commitDeadlineSeconds: especifica o prazo em segundos para a chamada de API Commit.
- bootstrapServer: servidor de inicialização do Kafka. Por exemplo,
localhost:9092. - kafkaTopic: tópico do Kafka em que a gravação será feita. Por exemplo,
topic.
Executar o modelo
Console
- Acesse a página Criar job usando um modelo do Dataflow. Acesse Criar job usando um modelo
- No campo Nome do job, insira um nome exclusivo.
- Opcional: em Endpoint regional, selecione um valor no menu suspenso. A região padrão é
us-central1.Para ver uma lista de regiões em que é possível executar um job do Dataflow, consulte Locais do Dataflow.
- No menu suspenso Modelo do Dataflow, selecione o modelo Gerador de dados de streaming.
- Nos campos de parâmetro fornecidos, insira os valores de parâmetro.
- Cliquem em Executar job.
gcloud
No shell ou no terminal, execute o modelo:
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
Substitua:
PROJECT_ID: o ID do projeto Google Cloud em que você quer executar o job do DataflowREGION_NAME: a região em que você quer implantar o job do Dataflow, por exemplo,us-central1JOB_NAME: um nome de job de sua escolhaVERSION: a versão do modelo que você quer usarUse estes valores:
latestpara usar a versão mais recente do modelo, disponível na pasta mãe não datada no bucket: gs://dataflow-templates-REGION_NAME/latest/- o nome da versão, como
2023-09-12-00_RC00, para usar uma versão específica do modelo, que pode ser encontrada aninhada na respectiva pasta mãe datada no bucket: gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: o caminho para o arquivo de esquema no Cloud Storage. Por exemplo,gs://mybucket/filename.json.QPS: o número de mensagens a serem publicadas por segundoPUBSUB_TOPIC: o tópico de saída do Pub/Sub. Por exemplo,projects/my-project-id/topics/my-topic-id.
API
Para executar o modelo usando a API REST, envie uma solicitação HTTP POST. Para mais informações sobre a
API e os respectivos escopos de autorização, consulte
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/", } }
Substitua:
PROJECT_ID: o ID do projeto Google Cloud em que você quer executar o job do DataflowLOCATION: a região em que você quer implantar o job do Dataflow, por exemplo,us-central1JOB_NAME: um nome de job de sua escolhaVERSION: a versão do modelo que você quer usarUse estes valores:
latestpara usar a versão mais recente do modelo, disponível na pasta mãe não datada no bucket: gs://dataflow-templates-REGION_NAME/latest/- o nome da versão, como
2023-09-12-00_RC00, para usar uma versão específica do modelo, que pode ser encontrada aninhada na respectiva pasta mãe datada no bucket: gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: o caminho para o arquivo de esquema no Cloud Storage. Por exemplo,gs://mybucket/filename.json.QPS: o número de mensagens a serem publicadas por segundoPUBSUB_TOPIC: o tópico de saída do Pub/Sub. Por exemplo,projects/my-project-id/topics/my-topic-id.
A seguir
- Saiba mais sobre os modelos do Dataflow.
- Confira a lista de modelos fornecidos pelo Google.