【发布时间】:2021-09-11 02:57:23
【问题描述】:
我已经使用 on_success_callback 和 on_failure_callback 实现了成功和失败的电子邮件警报。
上下文字典作为单个参数传递给此函数。
如何将另一个参数传递给这些回调方法?
这是我的代码
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