Airflow中如何获取SQL执行值并基于结果触发邮件任务?
解决方案
步骤1:修改run_sql.py输出结果到标准输出
BashOperator默认只会将bash命令的*标准输出(stdout)*推送到XCom,你之前拿不到值是因为run_sql.py里的res变量没有输出到stdout,且BashOperator没有开启推送开关。
在run_sql.py的末尾添加打印逻辑,仅输出res的值到标准输出,其他无关日志可以调整输出到stderr避免干扰:
# run_sql.py 原有执行逻辑结束后添加 if __name__ == "__main__": # 原有逻辑中res变量存储SQL执行结果 print(res) # 仅打印res值,不要添加其他额外打印内容
步骤2:调整BashOperator配置开启XCom推送
给你的BashOperator添加do_xcom_push=True参数,同时修正原有命令的语法问题:
check_phase = BashOperator( task_id='check_phase', bash_command='python3 /path/to/run_sql.py -e dev /path/to/sql/run_sql.sql', do_xcom_push=True, # 开启后会把bash命令的stdout推送到XCom dag=dag, )
步骤3:下游任务获取XCom值做判断
下游任务可以通过ti.xcom_pull()方法拿到check_phase任务输出的结果,注意拿到的结果是字符串类型,需要转换为int再做判断,示例如下:
from airflow.operators.python import BranchPythonOperator from airflow.operators.email import EmailOperator def check_res(ti): # 拉取XCom值,指定task_id为上游的check_phase sql_res = ti.xcom_pull(task_ids='check_phase') # 先去掉输出前后的换行、空白字符,再转成整数判断 res_int = int(sql_res.strip()) if res_int == 0: # 返回等于0时要执行的邮件任务ID return "send_success_email" else: # 返回不等于0时要执行的邮件任务ID return "send_fail_email" # 分支判断任务 check_res_task = BranchPythonOperator( task_id='check_res_task', python_callable=check_res, dag=dag ) # 邮件发送任务示例 send_success_email = EmailOperator( task_id='send_success_email', to='你的邮箱地址', subject='SQL执行结果为0,运行成功', html_content='<p>SQL任务执行正常,返回值为0</p>', dag=dag ) send_fail_email = EmailOperator( task_id='send_fail_email', to='你的邮箱地址', subject='SQL执行结果非0,运行异常', html_content='<p>SQL任务执行异常,返回值不为0</p>', dag=dag ) # 配置任务依赖关系 check_phase >> check_res_task >> [send_success_email, send_fail_email]
注意事项
- 确保run_sql.py中除了
print(res)之外,没有其他打印到stdout的内容,否则会导致XCom拉取的结果包含额外字符,转换int时报错 - 如果不需要分支执行,也可以直接在同一个PythonOperator里完成结果判断+邮件发送的逻辑
内容的提问来源于stack exchange,提问作者Sudarshan
相关产品推荐
相关产品推荐

