使用Apache Airflow進行編排

Apache Airflow是一個開源平台,可讓您建立、安排和監控工作流程。 工作流程被定義為有向無環圖 (DAG),它由使用Python程式碼編寫的多個任務組成。 每個任務代表一個離散的工作單元,例如執行腳本、查詢資料庫或呼叫 API。 Airflow 架構支援擴充和並行執行,使其適合管理複雜的資料密集型管道。

Apache Airflow 支援以下用例:

  • ETL 或 ELT 管道:從各種來源提取數據,對其進行轉換,並將其載入到資料倉儲中。
  • 資料倉儲:在資料倉儲中安排定期更新和資料轉換。
  • 資料處理:跨不同系統協調分散式資料處理任務。

必要條件

  • Apache Airflow獨立活動實例。
  • watsonx.data的使用者 API 金鑰(使用者名稱和 api_key)。 例如,username: yourid@example.comapi_key: sfw....cv23
  • watsonx.data的 CRN (wxd_instance_id)。 從watsonx.data資訊頁面取得實例 ID。
  • 來自活動 Spark 引擎的 Spark 引擎 ID (spark_engine_id)。
  • 來自活動Presto引擎的Presto外部 url (presto_ext_url)。
  • 系統信任的 SSL 憑證位置(如果適用)。
  • 與 Spark 和Presto引擎關聯的目錄 (catalog_name)。
  • 與所選目錄關聯的儲存桶的名稱。(儲存桶名稱)。
  • 使用以下指令安裝套件、Pandas 和Presto-python-client:pip install pandas presto-python-client

程序

  1. 此用例考慮將資料攝取到Presto任務。 為此,請建立一個 Spark 應用程序,將 Iceberg 資料提取到watsonx.data目錄。 這裡考慮範例Python檔ingestion-job.py。

    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. 將檔案上傳到名稱為 bucket_name 儲存。 有關更多信息,請參閱 將一些物件新增至您的儲存桶

  3. 使用Python設計 DAG 工作流程,並將Python檔案儲存到Apache Airflow目錄位置 $AIRFLOW_HOME/dags/ 目錄(AIRFLOW_HOME 的預設值設為 ~/airflow)。

    以下是一個工作流程範例,該工作流程執行任務以將資料擷取到watsonx.data中的Presto,並從watsonx.data查詢資料。 將文件儲存為以下內容,如 wxd_pipeline.py.

    使用 v2 API 時,請將 <api_version> 參數設定為 v2;使用 v3 API 時,請將它設定為 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. 登入 Apache Airflow

  5. 搜尋 wxd_pipeline.py job,從 Apache Airflow 控制台頁面啟用 DAG。 工作流程成功執行。