Dépendances de pipeline dans Data Fusion

Aug 28 2020

J'ai trois pipelines dans Data Fusion, à savoir A, B et C. Je veux que le pipeline C soit déclenché après l'exécution des pipelines A et B tous les deux terminés. Les déclencheurs de pipeline placent la dépendance sur un seul pipeline. Cela peut-il être implémenté dans Data Fusion?

Réponses

5 GonzaloPérezFernández Aug 28 2020 at 10:56

Vous pouvez le faire à l'aide de Google Cloud Composer [1]. Pour effectuer cette action, vous devez tout d'abord créer un nouvel environnement dans Google Cloud Composer [2], une fois terminé, vous devez installer un nouveau package Python dans votre environnement [3] et le package dont vous aurez besoin l'installation est [4] "apache-airflow-backport-provider-google".

Avec ce package installé, vous pourrez utiliser ces opérations [5], celle dont vous aurez besoin est [6] "Démarrer un pipeline DataFusion", de cette façon vous pourrez démarrer un nouveau pipeline depuis Airflow.

Un exemple de code python serait le suivant:

import airflow
import datetime
from airflow import DAG
from airflow import models
from airflow.operators.bash_operator import BashOperator
from datetime import timedelta
from airflow.providers.google.cloud.operators.datafusion import (
    CloudDataFusionStartPipelineOperator
)

default_args = {
   'start_date': airflow.utils.dates.days_ago(0),
   'retries': 1,
   'retry_delay': timedelta(minutes=5)
}

with models.DAG(
    'composer_DF',
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_args) as dag:

    # the operations.
    A = CloudDataFusionStartPipelineOperator(
            location="us-west1", pipeline_name="A", 
            instance_name="instance_name", task_id="start_pipelineA",
        )
    B = CloudDataFusionStartPipelineOperator(
            location="us-west1", pipeline_name="B", 
            instance_name="instance_name", task_id="start_pipelineB",
        )
    C = CloudDataFusionStartPipelineOperator(
            location="us-west1", pipeline_name="C", 
            instance_name="instance_name", task_id="start_pipelineC",
        )
    # First A then B and then C
    A >> B >> C

Vous pouvez définir les intervalles de temps en consultant la documentation Airflow.

Une fois ce code enregistré en tant que fichier .py, enregistrez-le dans le dossier Google Cloud Storage DAG de votre environnement.

Lorsque le DAG démarre, il exécute la tâche A et une fois terminé, il exécute la tâche B et ainsi de suite.

[1] https://cloud.google.com/composer

[2] https://cloud.google.com/composer/docs/how-to/managing/creating#:~:text=In%20the%20Cloud%20Console%2C%20open%20the%20Create%20Environment%20page.&text=Under%20Node%20configuration%2C%20click%20Add%20environment%20variable.&text=The%20From%3A%20email%20address%2C%20such,%40%20.&text=Your%20SendGrid%20API%20key.

[3] https://cloud.google.com/composer/docs/how-to/using/installing-python-dependencies

[4] https://pypi.org/project/apache-airflow-backport-providers-google/

[5] https://airflow.readthedocs.io/en/latest/_api/airflow/providers/google/cloud/operators/datafusion/index.html

[6] https://airflow.readthedocs.io/en/latest/howto/operator/google/cloud/datafusion.html#start-a-datafusion-pipeline

1 narendra Sep 08 2020 at 16:35

Il n'y a pas de moyen direct auquel je pourrais penser, mais deux solutions de contournement

Travaillez autour de 1 . La fusion des pipelines A et B dans le pipeline AB puis déclenche le pipeline C (AB> C).

Pipeline A - (GCS Copy> Decompress), Pipeline B - (GCS2> thrashsad)

BigQueryExecute pour atténuer l'erreur: DAG non valide. Il y a une île composée d'étapes.

Dans BigQueryExecute, requête valide et factice.

La fusion des deux pipelines en un seul peut compliquer les tests de pipeline. Pour surmonter cela, vous pouvez ajouter une condition fictive pour exécuter un pipeline une fois.

  1. Dans BigQueryExecute, remplacez la requête par «Select $ {flag}» et transmettez la valeur de l'indicateur dans l'argument d'exécution ou Sélectionnez 1 comme indicateur et cochez «Row As Arguments» sur true.
  2. Ajouter un plug-in de condition après BigQueryExecute et mettre la condition d'exécution ['flag'] = 1
  3. Le plugin Condition a deux sorties, connectez-les au pipeline A et au pipeline B.

Solution 2 : stockez l'indicateur des deux pipelines (A et B) dans la table BiqQuery, créez deux flux A> C et B> C pour déclencher le pipeline C. Cela déclencherait le pipeline C deux fois, mais l'utilisation de BigQueryExecute et le plug-in de condition ne s'exécuteront que lorsque les deux indicateurs sont disponibles dans la table BigQuery.

Comment?

  1. Dans les pipelines A et B pour écrire la sortie (une ligne) dans la table BigQuery "Pipeline_Run"
  2. Dans Pipeline C, ajoutez BigQueryExecute et interrogez "select count (*) as Cnt from ds.Pipeline_Run" et cochez "Row As Arguments" sur true.
  3. Dans Pipeline C, ajoutez le plugin Condition et vérifiez si la valeur de cnt est 2 (runtime ['cnt'] = 2) et connectez le reste des plugins du pipeline à sa sortie "Yes".
adiideas Sep 02 2020 at 11:20

Vous pouvez explorer les «horaires» définis via les API REST CDAP. Cela permet une exécution parallèle des pipelines et il n'y a aucune dépendance sur Cloud Composer (sauf pour le déclencheur basé sur un fichier du premier pipeline dans le workflow. Pour cela, vous auriez besoin d'une fonction cloud ou vous pourriez être un capteur de fichier Cloud composer)