Modelo do gerador de dados de streaming

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

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:

  1. 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.

  2. Faça upload do arquivo de esquema para um bucket do Cloud Storage.
  3. Defina o parâmetro de modelo schemaLocation como 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:

  1. Defina o parâmetro de modelo outputType como AVRO para o formato Avro ou PARQUET para o formato Parquet.
  2. Crie um arquivo de esquema Avro.
  3. Faça upload do arquivo de esquema para o Cloud Storage.
  4. Defina o parâmetro de modelo avroSchemaLocation como 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

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 chamada OUTPUT_TABLE_error_records, em que OUTPUT_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ão Ns (segundos), Nm (minutos) e Nh (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 SQL INSERT INTO usada 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áusula VALUES. 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

  1. Acesse a página Criar job usando um modelo do Dataflow.
  2. Acesse Criar job usando um modelo
  3. No campo Nome do job, insira um nome exclusivo.
  4. 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.

  5. No menu suspenso Modelo do Dataflow, selecione o modelo Gerador de dados de streaming.
  6. Nos campos de parâmetro fornecidos, insira os valores de parâmetro.
  7. 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 Dataflow
  • REGION_NAME: a região em que você quer implantar o job do Dataflow, por exemplo, us-central1
  • JOB_NAME: um nome de job de sua escolha
  • VERSION: a versão do modelo que você quer usar

    Use estes valores:

  • 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 segundo
  • PUBSUB_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 Dataflow
  • LOCATION: a região em que você quer implantar o job do Dataflow, por exemplo, us-central1
  • JOB_NAME: um nome de job de sua escolha
  • VERSION: a versão do modelo que você quer usar

    Use estes valores:

  • 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 segundo
  • PUBSUB_TOPIC: o tópico de saída do Pub/Sub. Por exemplo, projects/my-project-id/topics/my-topic-id.

A seguir