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.comet '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
-
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]) -
Téléchargez le fichier vers le stockage avec le nom
bucket_name. Pour plus d'informations, voir Ajouter des objets à vos buckets. -
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 valeurv3.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 -
Connectez-vous à Apache Airflow.
-
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.