Le générateur de tâches vous permet de créer des tâches Dataflow par lot et en flux continu personnalisées. Vous pouvez également enregistrer des tâches du générateur de tâches en tant que fichiers Apache Beam YAML pour les partager et les réutiliser.
Créer un pipeline
Pour créer un pipeline dans le générateur de tâches, procédez comme suit :
Accédez à la page Tâches de la Google Cloud console.
Cliquez sur Créer une tâche à partir du générateur.
Dans le champ Nom de la tâche, saisissez un nom pour la tâche.
Sélectionnez Lot ou Flux continu.
Si vous sélectionnez Flux continu, sélectionnez un mode de fenêtrage. Saisissez ensuite une spécification pour la fenêtre, comme suit :
- Fenêtre fixe : saisissez une taille de fenêtre en secondes.
- Fenêtre glissante : saisissez une taille et une durée de fenêtre, en secondes.
- Fenêtre de session : saisissez un intervalle de session en secondes.
Pour en savoir plus sur le fenêtrage, consultez la page Windows et fonctions de fenêtrage.
Ensuite, ajoutez des sources, des transformations et des récepteurs au pipeline, comme décrit dans les sections suivantes.
Ajouter une source au pipeline
Un pipeline doit comporter au moins une source. Initialement, le générateur de tâches est renseigné avec une source vide. Pour configurer la source, procédez comme suit :
Dans le champ Nom de la source, saisissez un nom pour la source ou utilisez le nom par défaut. Le nom apparaît dans le graphique de tâches lorsque vous exécutez la tâche.
Dans la liste Type de source, sélectionnez le type de source de données.
Selon le type de source, fournissez des informations de configuration supplémentaires.
- Par exemple, si vous sélectionnez BigQuery, spécifiez la table à partir de laquelle lire.
- Si vous sélectionnez Pub/Sub, spécifiez un schéma de message. Saisissez le nom et le type de données de chaque champ que vous souhaitez lire à partir des messages Pub/Sub. Le pipeline supprime tous les champs qui ne sont pas spécifiés dans le schéma.
- Si vous sélectionnez Apache Iceberg, spécifiez les informations de connexion de votre catalogue REST Iceberg (IRC), telles que l'identifiant de la table Iceberg, le nom du catalogue, le type de catalogue, l'URI du catalogue et le nom de l'entrepôt.
- Si vous sélectionnez Delta Lake, spécifiez les détails de votre table Delta Lake stockée dans Cloud Storage.
Facultatif : Pour certains types de sources, vous pouvez cliquer sur Prévisualiser les données sources pour prévisualiser les données sources.
Pour ajouter une autre source au pipeline, cliquez sur Ajouter une source. Pour combiner des données provenant de plusieurs sources, ajoutez une transformation SQL ou Join à votre pipeline.
Ajouter une transformation au pipeline
Vous pouvez éventuellement ajouter une ou plusieurs transformations au pipeline. Vous pouvez utiliser les transformations suivantes pour manipuler, agréger ou joindre des données provenant de sources et d'autres transformations :
| Type de transformation | Description | Informations sur la transformation Beam YAML |
|---|---|---|
| Filtrer (Python) | Filtrer les enregistrements avec une expression Python. | |
| Transformation SQL | Manipulez des enregistrements ou joignez plusieurs entrées avec une instruction SQL. | |
| Mapper les champs (Python) | Ajouter de nouveaux champs ou remapper des enregistrements entiers avec des expressions et des fonctions Python. | |
| Mapper les champs (SQL) | Ajouter ou mapper des champs d'enregistrement avec des expressions SQL. | |
Transformations YAML :
|
Utilisez n'importe quelle transformation du SDK Beam YAML. Configuration de la transformation YAML : indiquez les paramètres de configuration de la transformation YAML sous la forme d'un mappage YAML. Les paires clé/valeur sont utilisées pour remplir la section de configuration de la transformation Beam YAML obtenue. Pour connaître les paramètres de configuration compatibles pour chaque type de transformation, consultez la documentation sur la transformation Beam YAML. Exemples de paramètres de configuration : Combinegroup_by: combine: Rejoindretype: equalities: fields: |
|
| Journal | Enregistrer les enregistrements dans les journaux de nœud de calcul de la tâche. | |
| Grouper par |
Combiner des enregistrements avec des fonctions comme count() et
sum().
|
|
| Rejoindre | Joindre plusieurs entrées sur des champs égaux. | |
| Fractionner | Fractionner les enregistrements en aplatissant les champs du tableau. |
Pour ajouter une transformation, procédez comme suit :
Cliquez sur Ajouter une transformation.
Dans le champ Transformation, saisissez un nom pour la transformation ou utilisez le nom par défaut. Le nom apparaît dans le graphique de tâches lorsque vous exécutez la tâche.
Dans la liste Type de transformation, sélectionnez le type de transformation.
Selon le type de transformation, fournissez des informations de configuration supplémentaires. Par exemple, si vous sélectionnez Filtrer (Python), saisissez une expression Python à utiliser comme filtre.
Sélectionnez l'étape d'entrée de la transformation. L'étape d'entrée est la source ou la transformation dont la sortie fournit l'entrée de cette transformation.
Ajouter un récepteur au pipeline
Un pipeline doit comporter au moins un récepteur. Initialement, le générateur de tâches est renseigné avec un récepteur vide. Pour configurer le récepteur, procédez comme suit :
Dans le champ Nom du récepteur, saisissez un nom pour le récepteur ou utilisez le nom par défaut. Le nom apparaît dans le graphique de tâches lorsque vous exécutez la tâche.
Dans la liste Type de récepteur, sélectionnez le type de récepteur.
Selon le type de récepteur, fournissez des informations de configuration supplémentaires. Par exemple, si vous sélectionnez le récepteur BigQuery, sélectionnez la table BigQuery dans laquelle écrire.
Sélectionnez l'étape d'entrée du récepteur. L'étape d'entrée est la source ou la transformation dont la sortie fournit l'entrée de cette transformation.
Pour ajouter un autre récepteur au pipeline, cliquez sur Ajouter un récepteur.
Exécuter le pipeline
Pour exécuter un pipeline à partir du générateur de tâches, procédez comme suit :
Facultatif : Définissez les options de la tâche Dataflow. Pour développer la section des options Dataflow, cliquez sur la flèche d'expansion.
Cliquez sur Run Job (Exécuter la tâche). Le générateur de tâches accède au graphique de tâches pour la tâche envoyée. Vous pouvez utiliser le graphique de tâches pour surveiller l'état de la tâche.
Valider le pipeline avant de le lancer
Pour les pipelines avec une configuration complexe, tels que les filtres Python et les expressions SQL, il peut être utile de vérifier la configuration du pipeline pour détecter les erreurs de syntaxe avant de le lancer. Pour valider la syntaxe du pipeline, procédez comme suit :
- Cliquez sur Valider pour ouvrir Cloud Shell et démarrer le service de validation.
- Cliquez sur Commencer la validation.
- Si une erreur est détectée lors de la validation, un point d'exclamation rouge s'affiche.
- Corrigez les erreurs détectées et vérifiez les corrections en cliquant sur Valider. Si aucune erreur n'est détectée, une coche verte s'affiche.
Exécuter avec gcloud CLI
Vous pouvez également exécuter des pipelines Beam YAML à l'aide de gcloud CLI. Pour exécuter un pipeline du générateur de tâches avec gcloud CLI :
Cliquez sur Enregistrer YAML pour ouvrir la fenêtre Enregistrer YAML.
Effectuez l'une des actions suivantes :
- Pour enregistrer dans Cloud Storage, saisissez un chemin d'accès Cloud Storage, puis cliquez sur Enregistrer.
- Pour télécharger un fichier local, cliquez sur Télécharger.
Exécutez la commande suivante dans votre shell ou terminal :
gcloud dataflow yaml run my-job-builder-job --yaml-pipeline-file=YAML_FILE_PATHRemplacez
YAML_FILE_PATHpar le chemin d'accès à votre fichier YAML, localement ou dans Cloud Storage.
Étape suivante
- Utiliser l'interface de surveillance des jobs Dataflow.
- Enregistrer et charger des définitions de tâches YAML dans le générateur de tâches
- Apprenez-en plus sur YAML Beam.