Orchestration avec Apache Airflow

Apache Airflow est une plateforme open-source qui vous permet de créer, de planifier et de contrôler le flux de travail. Les flux de travail sont définis comme des graphes acycliques dirigés (DAG) qui consistent en de multiples tâches écrites à l'aide du code Python. Chaque tâche représente une unité de travail discrète, telle que l'exécution d'un script, l'interrogation d'une base de données ou l'appel d'une API. L'architecture Airflow prend en charge la mise à l'échelle et l'exécution parallèle, ce qui la rend adaptée à la gestion de pipelines complexes à forte intensité de données.

Apache airflow prend en charge les cas d'utilisation suivants :

  • Pipelines ETL ou ELT : Extraction de données à partir de diverses sources, transformation et chargement dans l'entrepôt de données.
  • Entreposage de données : Programmation de mises à jour régulières et de transformations de données dans un entrepôt de données.
  • Traitement des données : Orchestrer les tâches de traitement des données distribuées entre différents systèmes.

Prérequis

  • Instance active autonome d'Apache Airflow.
  • Clés d'API utilisateur pour watsonx.data (nom d'utilisateur et clé d'api). Par exemple, 'username: 'yourid@example.com et 'api_key: 'sfw....cv23.
  • CRN pour watsonx.data (wxd_instance_id). Obtenez l'ID de l'instance à partir de la page d'information watsonx.data.
  • Identifiant du moteur Spark d'un moteur Spark actif (spark_engine_id).
  • Presto url externe d'un moteur Presto actif (presto_ext_url).
  • Emplacement du certificat SSL auquel le système fait confiance (le cas échéant).
  • Catalogue associé aux moteurs Spark et Presto (catalog_name).
  • Nom du godet associé au catalogue sélectionné. (nom_du_seau).
  • Installez les paquets Pandas et Presto-python-client à l'aide de la commande : pip install pandas presto-python-client.

Procédure

  1. Le cas d'utilisation considère une tâche d'ingestion de données dans Presto. Pour ce faire, créez une application Spark qui ingère les données Iceberg dans le catalogue watsonx.data Ici, l'exemple de fichier Python ingestion-job.py est considéré.

    from pyspark.sql import SparkSession
    import os, sys
    
    def init_spark():
        spark = SparkSession.builder.appName("ingestion-demo").enableHiveSupport().getOrCreate()
        return spark
    
    def create_database(spark,bucket_name,catalog):
        spark.sql("create database if not exists {}.demodb LOCATION 's3a://{}/demodb'".format(catalog,bucket_name))
    
    def list_databases(spark,catalog):
        # list the database under lakehouse catalog
        spark.sql("show databases from {}".format(catalog)).show()
    
    def basic_iceberg_table_operations(spark,catalog):
        # demonstration: Create a basic Iceberg table, insert some data and then query table
        print("creating table")
        spark.sql("create table if not exists {}.demodb.testTable(id INTEGER, name VARCHAR(10), age INTEGER, salary DECIMAL(10, 2)) using iceberg".format(catalog)).show()
        print("table created")
        spark.sql("insert into {}.demodb.testTable values(1,'Alan',23,3400.00),(2,'Ben',30,5500.00),(3,'Chen',35,6500.00)".format(catalog))
        print("data inserted")
        spark.sql("select * from {}.demodb.testTable".format(catalog)).show()
    
    
    
    def clean_database(spark,catalog):
        # clean-up the demo database
        spark.sql("drop table if exists {}.demodb.testTable purge".format(catalog))
        spark.sql("drop database if exists {}.demodb cascade".format(catalog))
    
    def main(wxdDataBucket, wxdDataCatalog):
        try:
            spark = init_spark()
    
            create_database(spark,wxdDataBucket,wxdDataCatalog)
            list_databases(spark,wxdDataCatalog)
            basic_iceberg_table_operations(spark,wxdDataCatalog)
    
    
        finally:
            # clean-up the demo database
            clean_database(spark,wxdDataCatalog)
            spark.stop()
    
    if __name__ == '__main__':
        main(sys.argv[1],sys.argv[2])
    
    
  2. Téléchargez le fichier vers le stockage avec le nom bucket_name. Pour plus d'informations, voir Ajouter des objets à vos buckets.

  3. Concevoir un flux de travail DAG à l'aide de Python et enregistrer le fichier Python dans le répertoire Apache Airflow, répertoire '$AIRFLOW_HOME/dags/ (la valeur par défaut de AIRFLOW_HOME est fixée à ~/airflow).

    Voici un exemple de flux de travail, qui exécute des tâches pour ingérer des données dans Presto dans watsonx.data, et pour interroger des données à partir de watsonx.data. Enregistrez le fichier avec le contenu suivant, sous la forme wxd_pipeline.py.

    Lorsque vous utilisez l'API v2, attribuez au paramètre <api_version> la valeur v2; pour l'API v3, attribuez-lui la valeur v3.

    
    from datetime import timedelta, datetime
    from time import sleep
    import prestodb
    import pandas as pd
    import base64
    import os # type: ignore
    
    # The DAG object
    from airflow import DAG
    
    # Operators
    from airflow.operators.python_operator import PythonOperator # type: ignore
    import requests
    
    # Initializing the default arguments
    default_args = {
        'owner': 'IBM watsonx.data',
        'start_date': datetime(2024, 3, 4),
        'retries': 3,
        'retry_delay': timedelta(minutes=5),
        'wxd_endpoint': 'https://us-south.lakehouse.cloud.ibm.com', # Host endpoint
        'wxd_instance_id': 'crn:...::', # watsonx.data CRN
        'wxd_username': 'yourid@example.com', # your email id
        'wxd_api_key': 'sfw....cv23', # IBM IAM Api Key
        'spark_engine_id': 'spark6', # Spark Engine id
        'catalog_name': 'my_iceberg_catalog', # Catalog name where data will be ingestion
        'bucket_name': 'my-wxd-bucket', # Bucket name (not display name) associated with the above catalog
        'presto_eng_host': '2ce72...d59.cise...5s20.lakehouse.appdomain.cloud', # Presto engine hostname (without protocol and port)
        'presto_eng_port': 30912 # Presto engine port (in numbers only)
    }
    
    # Instantiate a DAG object
    wxd_pipeline_dag = DAG('wxd_ingestion_pipeline_saas',
            default_args=default_args,
            description='watsonx.data ingestion pipeline',
            schedule_interval=None,
            is_paused_upon_creation=True,
            catchup=False,
            max_active_runs=1,
            tags=['wxd', 'watsonx.data']
    )
    
    # Workaround: Enable if you want to disable SSL verification
    os.environ['NO_PROXY'] = '*'
    
    
    # Get access token
    def get_access_token():
        try:
            url = f"https://iam.cloud.ibm.com/oidc/token"
            headers = {
                'Content-Type': 'application/x-www-form-urlencoded',
                'Accept': 'application/json',
            }
    
            data = {
                'grant_type': 'urn:ibm:params:oauth:grant-type:apikey',
                'apikey': default_args['wxd_api_key'],
            }
    
    
            response = requests.post('https://iam.cloud.ibm.com/identity/token', headers=headers, data=data)
    
            return response.json()['access_token']
        except Exception as inst:
            print('Error in getting access token')
            print(inst)
            exit
    
    
    def _ingest_via_spark_engine():
        try:
            print('ingest__via_spark_engine')
            url = f"{default_args['wxd_endpoint']}/lakehouse/api/<api_version>/spark_engines/{default_args['spark_engine_id']}/applications"
    
            headers = {'Content-type': 'application/json', 'Authorization': f'Bearer {get_access_token()}', 'AuthInstanceId': default_args['wxd_instance_id']}
            auth_str = base64.b64encode(f'ibmlhapikey_{default_args["wxd_username"]}:{default_args["wxd_api_key"]}'.encode('ascii')).decode("ascii")
    
            response = requests.post(url, None, {
                "application_details": {
                    "conf": {
                        "spark.executor.cores": "1",
                        "spark.executor.memory": "1G",
                        "spark.driver.cores": "1",
                        "spark.driver.memory": "1G",
                        "spark.hadoop.wxd.apikey": f"Basic {auth_str}"
                    },
                    "application": f"s3a://{default_args['bucket_name']}/ingestion-job.py",
                    "arguments": [
                        default_args['bucket_name'],
                        default_args['catalog_name']
                    ],
                }
            } , headers=headers, verify=False)
    
            print("Response", response.content)
            return response.json()['id']
        except Exception as inst:
            print(inst)
            raise ValueError('Task failed due to', inst)
    
    
    def _wait_until_job_is_complete(**context):
        try:
            print('wait_until_job_is_complete')
            application_id = context['task_instance'].xcom_pull(task_ids='ingest_via_spark_engine')
            print(application_id)
    
            while True:
                url = f"{default_args['wxd_endpoint']}/lakehouse/api/<api_version>/spark_engines/{default_args['spark_engine_id']}/applications/{application_id}"
                headers = {'Content-type': 'application/json', 'Authorization': f'Bearer {get_access_token()}', 'AuthInstanceId': default_args['wxd_instance_id']}
    
                response = requests.get(url, headers=headers, verify=False)
                print(response.content)
    
                data = response.json()
    
                if data['state'] == 'finished':
                    break
                elif data['state'] in ['stopped', 'failed', 'killed']:
                    raise ValueError("Job failed: ", data)
    
                print('Job is not completed, sleeping for 10secs')
                sleep(10)
        except Exception as inst:
            print(inst)
            raise ValueError('Task failed due to', inst)
    
    
    def _query_presto():
        try:
            with prestodb.dbapi.connect(
                host=default_args['presto_eng_host'],
                port=default_args['presto_eng_port'],
                user=default_args['wxd_username'],
                catalog='tpch',
                schema='tiny',
                http_scheme='https',
                auth=prestodb.auth.BasicAuthentication(f'ibmlhapikey_{default_args["wxd_username"]}', default_args["wxd_api_key"])
            ) as conn:
                df = pd.read_sql_query(f"select * from {default_args['catalog_name']}.demodb.testTable limit 5", conn)
    
                with pd.option_context('display.max_rows', None, 'display.max_columns', None):
                    print("\n", df.head())
        except Exception as inst:
            print(inst)
            raise ValueError('Query faield due to ', inst)
    
    
    def start_job():
        print('Validating default arguments')
    
        if 'wxd_endpoint' not in default_args:
            raise ValueError('wxd_endpoint is mandatory')
    
        if 'wxd_username' not in default_args:
            raise ValueError('wxd_username is mandatory')
    
        if 'wxd_instance_id' not in default_args:
            raise ValueError('wxd_instance_id is mandatory')
    
        if 'wxd_api_key' not in default_args:
            raise ValueError('wxd_api_key is mandatory')
    
        if 'spark_engine_id' not in default_args:
            raise ValueError('spark_engine_id is mandatory')
    
    
    start = PythonOperator(task_id='start_task', python_callable=start_job, dag=wxd_pipeline_dag)
    
    ingest_via_spark_engine = PythonOperator(task_id='ingest_via_spark_engine', python_callable=_ingest_via_spark_engine, dag=wxd_pipeline_dag)
    wait_until_ingestion_is_complete = PythonOperator(task_id='wait_until_ingestion_is_complete', python_callable=_wait_until_job_is_complete, dag=wxd_pipeline_dag)
    query_via_presto = PythonOperator(task_id='query_via_presto', python_callable=_query_presto, dag=wxd_pipeline_dag)
    
    start >> ingest_via_spark_engine >> wait_until_ingestion_is_complete >> query_via_presto
    
    
  4. Connectez-vous à Apache Airflow.

  5. Recherchez le travail wxd_pipeline.py, activez le DAG à partir de la page de la console Apache Airflow. Le flux de travail est exécuté avec succès.