【问题标题】:How to pass parameters to Airflow on_success_callback and on_failure_callback如何将参数传递给 Airflow on_success_callback 和 on_failure_callback
【发布时间】:2021-09-11 02:57:23
【问题描述】:

我已经使用 on_success_callback 和 on_failure_callback 实现了成功和失败的电子邮件警报。

根据Airflow documentation

上下文字典作为单个参数传递给此函数。

如何将另一个参数传递给这些回调方法?

这是我的代码

from airflow.utils.email import send_email_smtp

def task_success_alert(context):
    subject = "[Airflow] DAG {0} - Task {1}: Success".format(
        context['task_instance_key_str'].split('__')[0], 
        context['task_instance_key_str'].split('__')[1]
        )
    html_content = """
    DAG: {0}<br>
    Task: {1}<br>
    Succeeded on: {2}
    """.format(
        context['task_instance_key_str'].split('__')[0], 
        context['task_instance_key_str'].split('__')[1], 
        datetime.now()
        )
    send_email_smtp(dag_vars["dev_mailing_list"], subject, html_content)

def task_failure_alert(context):
    subject = "[Airflow] DAG {0} - Task {1}: Failed".format(
        context['task_instance_key_str'].split('__')[0], 
        context['task_instance_key_str'].split('__')[1]
        )
    html_content = """
    DAG: {0}<br>
    Task: {1}<br>
    Failed on: {2}
    """.format(
        context['task_instance_key_str'].split('__')[0], 
        context['task_instance_key_str'].split('__')[1], 
        datetime.now()
        )
    send_email_smtp(dag_vars["dev_mailing_list"], subject, html_content)

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2019, 6, 13),
    'on_success_callback': task_success_alert,
    'on_failure_callback': task_failure_alert
}

我打算将回调移动到另一个包并将电子邮件地址作为参数传递。

【问题讨论】:

  • 当我使用 CONTEXT 时,我的任务完成了,但 dag 继续运行。它永远不会结束。有什么建议么。我正在使用 def task_alert(context): dag_id = context['dag'].dag_id \n task_id = context['task_instance']。 task_id 我正在调用 on_failure_callback=task_alert
  • 使用 xcom_push 和 xcom_pull

标签: python airflow airflow-scheduler


【解决方案1】:

您可以在 dag 中定义一个函数,该函数从您的包中调用该函数。在调用该函数时,将电子邮件作为参数传递。您可以在 DAG 级别进一步细化它,以仅传递电子邮件所需的信息。

from package import outer_task_success_callback
email = 'xyz@example.com'

def task_success_alert(context):
    dag_id = context['dag'].dag_id
    task_id = context['task_instance']. task_id
    outer_task_success_callback(dag_id, task_id, email)
    
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2019, 6, 13),
    'on_success_callback': task_success_alert,
    'on_failure_callback': task_failure_alert
}

这将允许您在调用包中的函数之前进行自定义。

附带说明,airflow 具有 smtp 电子邮件功能。您可以利用这些解决方案,而不是编写自己的解决方案。

【讨论】:

  • 当我使用 CONTEXT 时,我的任务完成了,但 dag 继续运行。它永远不会结束。任何建议。
  • 没有代码示例无法理解或调试。添加一个问题并将其标记为气流。如果可以的话,我会检查并回答
  • def notify_email(context): import inspect """Send custom email alerts.""" import smtplib, ssl from email.mime.text import MIMEText from email.mime.multipart import MIMEMultipart sender_email = ' abc@gmail.com' receiver_email = 'xyz@gmail.com' 密码 = "abc" message = MIMEMultipart("alternative") #task_instance = context['task'].task_id dag_instance=context['dag_id'].dag_id 什么时候使用 context['dag_id'].dag_id 我的 dag 继续运行任务完成并且邮件没有发送
  • @Gabriel Eckers 气流中的 SMTP 实用程序仍然提供将电子邮件作为任务发送的选项。虽然 on_success 选项不可用,但电子邮件功能可以保留在气流 DAG 中。
【解决方案2】:

您可以使用 partial 创建一个带有预定义参数的函数,例如:

from functools import partial
new_task_success_alert = partial(task_success_alert, email='your_email')

然后将新函数添加为回调:

on_success_callback=new_task_success_alert

【讨论】:

    【解决方案3】:

    您可以创建一个任务,其唯一目的是通过xcoms 推送配置设置。您可以通过context 拉取配置,因为task_instance 对象包含在context 中。

    def push_configuration(ti, params):
        ti.xcom_push(key='conn_id', value=params)
    
    def task_success_alert(context):
        ti = context.get('ti') 
        params = ti.xcom_pull(key='params', task_ids='Settings')
        ...
    
    
    step0 = PythonOperator(
            task_id='Settings',
            python_callable=push_configuration,
            op_kwargs={'params': params})
    
    step1 = BashOperator(
            task_id='step1',
            bash_command='pwd',
            on_success_callback=task_success_alert)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-02-09
      • 2021-10-24
      • 1970-01-01
      • 1970-01-01
      • 2018-05-03
      • 1970-01-01
      • 2020-07-11
      相关资源
      最近更新 更多