【问题标题】:Airflow on_success_callback and on_failure_callback not working with data bricks notebook气流 on_success_callback 和 on_failure_callback 不适用于数据砖笔记本
【发布时间】:2021-01-12 13:46:30
【问题描述】:

我想自定义我的 DAG 以在成功或失败时调用 datarbicks 笔记本。我创建了两个不同的函数来根据成功/失败案例调用 databricks 笔记本。成功或失败回调函数正在调用,但 databricks 笔记本未执行。这是示例代码。

def task_success_callback(context):
    """ task_success callback """
    context['task_instance'].task_id
    print("success case")
    dq_notebook_success_task_params = {
        'existing_cluster_id': Variable.get("DATABRICKS_CLUSTER_ID"),
        'notebook_task': {
            'notebook_path': '/AAA/Airflow/Operators/audit_file_operator',
             'base_parameters': {
                "root": "dbfs:/mnt/aaa",
                "audit_file_path": "/success_file_path/",
                "table_name": "sample_data_table",
                "audit_flag": "success"
            }
        }
    }

    DatabricksSubmitRunOperator(
    task_id="weather_table_task_id",
    databricks_conn_id='databricks_conn',
    json=dq_notebook_success_task_params,
    do_xcom_push=True,
    secrets=[secret.Secret(
    deploy_type='env',
    deploy_target=None,
    secret='adf-service-principal'
    ), secret.Secret(
    deploy_type='env',
    deploy_target=None,
    secret='postgres-credentials',
    )],
    )

def task_failure_callback(context):
    """ task_success callback """
    context['task_instance'].task_id
    print("failure case")
    dq_notebook_failure_task_params = {
        'existing_cluster_id': Variable.get("DATABRICKS_CLUSTER_ID"),
        'notebook_task': {
            'notebook_path': '/AAA/Airflow/Operators/audit_file_operator',
            'base_parameters': {
                "root": "dbfs:/mnt/aaa",
                "audit_file_path": "/failure_file_path/",
                "table_name": "sample_data_table",
                "audit_flag": "failure"
            }
        }
    }

    DatabricksSubmitRunOperator(
    task_id="weather_table_task_id",
    databricks_conn_id='databricks_conn',
    json=dq_notebook_failure_task_params,
    do_xcom_push=True,
    secrets=[secret.Secret(
    deploy_type='env',
    deploy_target=None,
    secret='adf-service-principal'
    ), secret.Secret(
    deploy_type='env',
    deploy_target=None,
    secret='postgres-credentials',
    )],
    )

DEFAULT_ARGS = {
    "owner": "admin",
    "depends_on_past": False,
    "start_date": datetime(2020, 9, 23),
    "on_success_callback": task_success_callback,
    "on_failure_callback": task_failure_callback,
    "email": ["airflow@airflow.com"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(seconds=10),
}

==================
Remaining DAG code
==================

【问题讨论】:

    标签: airflow databricks


    【解决方案1】:

    在 Airflow 中,每个运算符都有 execute() 方法来定义运算符逻辑。当您创建工作流时,气流初始化构造函数,渲染模板并为您调用执行方法。但是,当您在 python 函数中定义运算符时,您还需要自己处理。

    所以当你写的时候:

    def task_success_callback(context):
       DatabricksSubmitRunOperator(..)
    

    您在这里所做的只是初始化DatabricksSubmitRunOperator 接触器。你没有调用操作符逻辑。

    你需要做的是:

    def task_success_callback(context):
       op = DatabricksSubmitRunOperator(..)
       op.execute()
    

    【讨论】:

    • 感谢 Elad 的回复,它对我有用。
    • @Kiran 如果解决了请接受答案:)
    • Elad,知道如何将参数传递给 task_success_callback(context) 函数吗?
    • Elad,感谢您的回复,但这仅适用于静态参数,但我需要将这些参数作为动态参数传递。你能帮我如何传递动态参数吗?
    【解决方案2】:
    TableList = collections.namedtuple(
        "table_list",
        "table_name audit_file_name",
    )
    LIST_OF_TABLES = [
        TableList(
            table_name="table1",
            audit_file_name="/testdata/Audit_files/",
        ),
        TableList(
            table_name="table2",
            audit_file_name="/testdata/Audit_files/",
        ),
        TableList(
            table_name="table3",
            audit_file_name="/testdata/Audit_files/",
        ),
        TableList(
            table_name="table4",
            audit_file_name="/testdata/Audit_files/",
        )
    ]
    for table in LIST_OF_TABLES:
        DEFAULT_ARGS = {
            "owner": "admin",
            "depends_on_past": False,
            "start_date": datetime(2020, 9, 23),
            "on_success_callback": partial(task_success_callback,table.table_name,table.audit_file_name),
            "on_failure_callback": partial(task_failure_callback,table.table_name,table.audit_file_name),
            "email": ["airflow@airflow.com"],
            "email_on_failure": False,
            "email_on_retry": False,
            "retries": 1,
            "retry_delay": timedelta(seconds=10),
        }
        WORKFLOW = DAG(
            'test_dag',
            default_args=DEFAULT_ARGS,
            schedule_interval="30 3 * * 1",
            catchup=False,
        )
    

    【讨论】:

    • Elad,这里我想为多个表加载的回调函数传递表名和 audit_file_name 路径。
    • @Elad,你对此有什么意见吗?
    猜你喜欢
    • 1970-01-01
    • 2021-05-04
    • 1970-01-01
    • 2017-12-31
    • 1970-01-01
    • 1970-01-01
    • 2017-11-18
    • 1970-01-01
    • 2022-01-26
    相关资源
    最近更新 更多