Modèle de générateur de flux de données

Le modèle de générateur de flux de données génère des enregistrements ou des messages synthétiques et les envoie à un récepteur de destination. Vous pouvez configurer le schéma d'enregistrement et la fréquence à laquelle les enregistrements sont générés.

Le modèle est compatible avec les récepteurs suivants :

  • Sujet Apache Kafka
  • Table BigQuery
  • Bucket Cloud Storage
  • Point de terminaison Java Database Connectivity (JDBC)
  • Sujet Pub/Sub
  • Table Spanner

Voici quelques cas d'utilisation possibles :

  • Simuler un événement en temps réel et à grande échelle sur un sujet Pub/Sub, afin de mesurer et de déterminer le nombre et la taille des consommateurs requis pour traiter les événements publiés
  • Générez des données synthétiques pour évaluer les benchmarks de performances ou pour servir de démonstration de faisabilité.
  • Validez un pipeline de bout en bout. Par exemple, envoyez des enregistrements à un sujet Kafka, qui sont ensuite lus par un consommateur en aval.

Conditions requises pour ce pipeline

Définir le schéma d'enregistrement

Le modèle fournit un schéma prédéfini pour les données générées. Pour utiliser ce schéma, définissez le paramètre de modèle schemaTemplate sur GAME_EVENT.

Vous pouvez également fournir votre propre schéma de données comme suit :

  1. Créez un fichier de schéma contenant un modèle JSON pour les données générées. Ce modèle utilise la bibliothèque JSON Data Generator, qui accepte différentes fonctions pour randomiser les données. Exemple :

    {
      "id": {{integer(0,1000)}},
      "name": "{{uuid()}}",
      "isInStock": {{bool()}}
    }

    Pour en savoir plus, consultez la documentation sur json-data-generator.

  2. Importez le fichier de schéma dans un bucket Cloud Storage.
  3. Définissez le paramètre de modèle schemaLocation sur l'URI Cloud Storage du fichier de modèle.

Spécifier le format de sortie

Par défaut, le modèle génère des données JSON. Pour certaines destinations, le modèle est également compatible avec les formats Avro ou Parquet :

  • Avro : compatible avec Cloud Storage, Apache Kafka et Pub/Sub
  • Parquet : compatible avec Cloud Storage.

Pour générer un format Avro ou Parquet :

  1. Définissez le paramètre de modèle outputType sur AVRO pour le format Avro ou sur PARQUET pour le format Parquet.
  2. Créez un fichier de schéma Avro.
  3. Importez le fichier de schéma dans Cloud Storage.
  4. Définissez le paramètre de modèle avroSchemaLocation sur l'URI Cloud Storage du fichier de schéma.

Spécifier le récepteur de destination

Les sections suivantes décrivent comment configurer le modèle pour chaque type de destination.

Sujet Apache Kafka

Pour écrire dans un sujet Kafka, définissez les paramètres de modèle suivants :

  • sinkType : KAFKA.
  • bootstrapServer : adresse d'amorçage du cluster Kafka.
  • kafkaTopic : sujet Kafka dans lequel écrire.

Si vous écrivez dans un cluster Google Cloud Managed Service pour Apache Kafka, accordez le rôle Client Managed Kafka (roles/managedkafka.client) au compte de service du worker.

Table BigQuery

Pour écrire dans une table BigQuery, définissez les paramètres de modèle suivants :

  • sinkType : BIGQUERY.
  • outputTableSpec : table BigQuery dans laquelle écrire. Mettez en forme ce paramètre comme suit : PROJECT_ID:DATASET.TABLE.

Les paramètres ci-dessous sont facultatifs :

  • outputDeadletterTable : nom de la table dans laquelle le pipeline écrit les enregistrements ayant échoué. Si cette valeur n'est pas spécifiée, le pipeline crée une table nommée OUTPUT_TABLE_error_records, où OUTPUT_TABLE correspond au nom de la table de sortie.
  • writeDisposition : spécifie comment écrire dans une table existante. Les valeurs suivantes sont acceptées :

    • WRITE_APPEND : ajoute des lignes à la table existante.
    • WRITE_TRUNCATE. Tronquez les lignes existantes.
    • WRITE_EMPTY. Écrivez uniquement si la table est vide. Si la table contient déjà des données, la tâche échoue.

    La valeur par défaut est WRITE_APPEND.

Attribuez le rôle Éditeur de données BigQuery (roles/bigquery.dataEditor) au compte de service du worker.

Cloud Storage

Pour écrire dans un bucket Cloud Storage, définissez les paramètres de modèle suivants :

  • sinkType : GCS.
  • outputDirectory : chemin d'accès au dossier Cloud Storage dans lequel écrire.

Les paramètres ci-dessous sont facultatifs :

  • numShards : nombre maximal de partitions. Une valeur plus élevée peut permettre un débit plus élevé, mais potentiellement des coûts d'agrégation de données plus élevés. Si la valeur est définie sur 0, Dataflow sélectionne le nombre de segments. La valeur par défaut est 0.
  • outputFilenamePrefix : préfixe du nom de fichier. La valeur par défaut est output-.
  • windowDuration : intervalle auquel le pipeline écrit des fichiers dans Cloud Storage. Les formats autorisés sont Ns (secondes), Nm (minutes) et Nh (heures). La valeur par défaut est 1m (1 minute).

Accordez au compte de service du worker le rôle Administrateur des objets de l'espace de stockage (roles/storage.objectAdmin).

Point de terminaison JDBC

Pour écrire dans un point de terminaison JDBC, définissez les paramètres de modèle suivants :

  • sinkType : JDBC.
  • driverClassName : classe du pilote JDBC à utiliser. Exemple : com.mysql.jdbc.Driver.
  • connectionUrl : chaîne de connexion à la source JDBC.
  • statement : instruction SQL INSERT INTO utilisée pour écrire dans la base de données. L'instruction doit spécifier les colonnes de table dans lesquelles écrire, avec des caractères d'espace réservé '?' pour la clause VALUES. Le pipeline remplace les espaces réservés par les valeurs de champ correspondantes issues des données JSON.

    Exemple : INSERT INTO tableName (column1, column2) VALUES (?,?).

Les paramètres ci-dessous sont facultatifs :

  • username : nom d'utilisateur pour la connexion JDBC.
  • password : mot de passe de la connexion JDBC
  • connectionProperties : chaîne de propriétés pour la connexion JDBC. Exemple : unicode=true;characterEncoding=UTF-8.

Sujet Pub/Sub

Pour écrire dans un sujet Pub/Sub, définissez les paramètres de modèle suivants :

  • sinkType : PUBSUB.
  • topic : sujet Pub/Sub dans lequel écrire.

Attribuez le rôle Diffuseur Pub/Sub (roles/pubsub.publisher) au compte de service du worker.

Table Spanner

Pour écrire dans une table Spanner, définissez les paramètres de modèle suivants :

  • sinkType : SPANNER.
  • projectId : ID du projet contenant la table Spanner.
  • spannerInstanceName : nom de l'instance Spanner.
  • spannerDatabaseName : nom de la base de données Spanner.
  • spannerTableName : nom de la table Spanner.

Les paramètres ci-dessous sont facultatifs :

  • maxNumMutations : nombre maximal de cellules mutées par lot.
  • maxNumRows : nombre maximal de lignes mutées par lot.
  • batchSizeBytes : nombre maximal d'octets mutés par lot.
  • commitDeadlineSeconds : délai pour l'appel d'API "commit", en secondes.

Attribuez le rôle Utilisateur de base de données Cloud Spanner (roles/spanner.databaseUser) au compte de service du worker.

Paramètres de modèle

Paramètres obligatoires

  • qps : indique le nombre de messages par seconde à publier dans Pub/Sub.

Paramètres facultatifs

  • schemaTemplate : modèle de schéma préexistant à utiliser. La valeur doit être l'une des suivantes : [GAME_EVENT].
  • schemaLocation : chemin d'accès Cloud Storage de l'emplacement du schéma. Exemple :gs://<bucket-name>/prefix
  • topic : nom du sujet dans lequel le pipeline doit publier des données. Exemple :projects/<project-id>/topics/<topic-name>
  • messagesLimit : indique le nombre maximal de messages de sortie à générer. 0 signifie illimité. La valeur par défaut est 0.
  • outputType : type de sortie du message. La valeur par défaut est JSON.
  • avroSchemaLocation : chemin d'accès Cloud Storage à l'emplacement du schéma Avro. Obligatoire lorsque le type de sortie est AVRO ou PARQUET. Exemple :gs://your-bucket/your-path/schema.avsc
  • sinkType : type de récepteur du message. La valeur par défaut est PUBSUB.
  • outputTableSpec : table BigQuery de sortie. Obligatoire lorsque sinkType est BIGQUERY. Par exemple, <project>:<dataset>.<table_name>.
  • writeDisposition : BigQuery WriteDisposition. Par exemple, WRITE_APPEND, WRITE_EMPTY ou WRITE_TRUNCATE. La valeur par défaut est WRITE_APPEND.
  • outputDeadletterTable : les messages n'ayant pas pu atteindre la table de sortie pour tous types de raisons (par exemple, schéma non concordant ou format JSON non valide) sont écrits dans cette table. Si elle n'existe pas, elle sera créée lors de l'exécution du pipeline. Exemple :your-project-id:your-dataset.your-table-name
  • windowDuration : durée/taille de la fenêtre dans laquelle les données seront écrites dans Cloud Storage. Les formats autorisés sont les suivants : Ns (pour les secondes, exemple : 5s), Nm (pour les minutes, exemple : 12m), Nh (pour les heures, exemple : 2h). Exemple :1m La valeur par défaut est "1m".
  • outputDirectory : chemin d'accès et préfixe du nom de fichier pour l'écriture des fichiers de sortie. Doit se terminer par une barre oblique. Le format DateTime permet d'analyser le chemin d'accès au répertoire pour les outils de mise en forme de date et d'heure. Exemple :gs://your-bucket/your-path/
  • outputFilenamePrefix : préfixe à placer sur chaque fichier ciblé sur une fenêtre. Exemple :output- La valeur par défaut est "output-".
  • numShards : nombre maximal de partitions de sortie générées lors de l'écriture. Un nombre plus élevé de segments entraîne un débit plus élevé pour l'écriture dans Cloud Storage, mais potentiellement un coût d'agrégation de données plus élevé entre les partitions lors du traitement des fichiers Cloud Storage de sortie. La valeur par défaut est déterminée par Dataflow.
  • driverClassName : nom de la classe du pilote JDBC à utiliser. Exemple :com.mysql.jdbc.Driver
  • connectionUrl : chaîne de connexion URL pour se connecter à la source JDBC. Exemple :jdbc:mysql://some-host:3306/sampledb
  • username : nom d'utilisateur à utiliser pour la connexion JDBC.
  • password : mot de passe à utiliser pour la connexion JDBC.
  • connectionProperties : chaîne de propriétés à utiliser pour la connexion JDBC. Le format de la chaîne doit être [propertyName=property;]*. Exemple :unicode=true;characterEncoding=UTF-8
  • statement : instruction SQL qui sera exécutée pour écrire dans la base de données. L'instruction doit spécifier les noms de colonnes de la table dans n'importe quel ordre. Seules les valeurs des noms de colonnes spécifiés seront lues à partir du fichier JSON et ajoutées à l'instruction. Exemple :INSERT INTO tableName (column1, column2) VALUES (?,?)
  • projectId : ID du projet GCP dans lequel se trouve la table Spanner.
  • spannerInstanceName : nom de l'instance Cloud Spanner.
  • spannerDatabaseName : nom de la base de données Cloud Spanner.
  • spannerTableName : nom de la table Cloud Spanner.
  • maxNumMutations : spécifie la limite de mutation des cellules (nombre maximal de cellules mutées par lot). La valeur par défaut est 5 000.
  • maxNumRows : spécifie la limite de mutation des lignes (nombre maximal de lignes mutées par lot). La valeur par défaut est 1 000.
  • batchSizeBytes : spécifie la limite de taille de lot (nombre maximal d'octets modifiés par lot). La valeur par défaut est de 1 Mo.
  • commitDeadlineSeconds : spécifie le délai en secondes pour l'appel d'API Commit.
  • bootstrapServer : serveur d'amorçage Kafka (par exemple, localhost:9092).
  • kafkaTopic : sujet Kafka dans lequel écrire. Exemple :topic

Exécuter le modèle

Console

  1. Accédez à la page Dataflow Créer un job à partir d'un modèle.
  2. Accéder à la page Créer un job à partir d'un modèle
  3. Dans le champ Nom du job, saisissez un nom de job unique.
  4. Facultatif : pour Point de terminaison régional, sélectionnez une valeur dans le menu déroulant. La région par défaut est us-central1.

    Pour obtenir la liste des régions dans lesquelles vous pouvez exécuter un job Dataflow, consultez la page Emplacements Dataflow.

  5. Dans le menu déroulant Modèle Dataflow, sélectionnez le modèle Générateur de données de flux.
  6. Dans les champs fournis, saisissez vos valeurs de paramètres.
  7. Cliquez sur Run Job (Exécuter la tâche).

gcloud

Dans le shell ou le terminal, exécutez le modèle :

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
  

Remplacez les éléments suivants :

  • PROJECT_ID : ID du projet Google Cloud dans lequel vous souhaitez exécuter le job Dataflow
  • REGION_NAME : région dans laquelle vous souhaitez déployer votre job Dataflow, par exemple us-central1
  • JOB_NAME : nom de job unique de votre choix
  • VERSION : version du modèle que vous souhaitez utiliser

    Vous pouvez utiliser les valeurs suivantes :

  • SCHEMA_LOCATION : chemin d'accès au fichier de schéma dans Cloud Storage. Exemple : gs://mybucket/filename.json.
  • QPS : nombre de messages à publier par seconde
  • PUBSUB_TOPIC : sujet Pub/Sub de sortie. Exemple : projects/my-project-id/topics/my-topic-id.

API

Pour exécuter le modèle à l'aide de l'API REST, envoyez une requête HTTP POST. Pour en savoir plus sur l'API, ses autorisations et leurs champs d'application, consultez la section 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/",
   }
}
  

Remplacez les éléments suivants :

  • PROJECT_ID : ID du projet Google Cloud dans lequel vous souhaitez exécuter le job Dataflow
  • LOCATION : région dans laquelle vous souhaitez déployer votre job Dataflow, par exemple us-central1
  • JOB_NAME : nom de job unique de votre choix
  • VERSION : version du modèle que vous souhaitez utiliser

    Vous pouvez utiliser les valeurs suivantes :

  • SCHEMA_LOCATION : chemin d'accès au fichier de schéma dans Cloud Storage. Exemple : gs://mybucket/filename.json.
  • QPS : nombre de messages à publier par seconde
  • PUBSUB_TOPIC : sujet Pub/Sub de sortie. Exemple : projects/my-project-id/topics/my-topic-id.

Étapes suivantes