如何在Airflow中根据数据库查询结果选择执行对应任务
基于SQL查询结果创建分支任务的实现方案
核心逻辑很明确:先执行指定SQL获取当日状态异常的记录数,再根据返回的count值分支执行对应任务——count=0触发成功报告任务,count=1触发失败报告任务。以下是几种常见场景的实现方式:
1. Apache Airflow 实现分支任务
Airflow的BranchPythonOperator专门用来处理分支逻辑,结合数据库Hook即可完成需求:
from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.operators.bash import BashOperator # 替换为你实际的任务算子 from airflow.providers.mysql.hooks.mysql import MySqlHook from datetime import datetime def decide_report_task(**context): # 获取任务启动日期,对应SQL中的{start_task_time} task_start_date = context['dag_run'].start_date.strftime('%Y-%m-%d') # 连接数据库执行查询 db_hook = MySqlHook(mysql_conn_id='your_mysql_connection_id') query = """ SELECT COUNT(*) FROM r_summary rs WHERE rs.status = 'Not Ok' AND rs.created_date = %s; """ abnormal_count = db_hook.get_first(query, parameters=(task_start_date,))[0] # 返回对应任务ID,触发分支 if abnormal_count == 1: return 'sending_reports_if_failed' elif abnormal_count == 0: return 'sending_reports_if_success' else: # 处理非预期的count值(比如多条异常记录) return 'handle_unexpected_abnormal_count' with DAG( 'daily_status_report_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily', catchup=False ) as dag: branch_decider = BranchPythonOperator( task_id='check_daily_status', python_callable=decide_report_task, provide_context=True ) # 定义实际任务,替换为你的业务逻辑 success_report = BashOperator( task_id='sending_reports_if_success', bash_command='python /path/to/your/success_report_script.py' ) failed_report = BashOperator( task_id='sending_reports_if_failed', bash_command='python /path/to/your/failed_report_script.py' ) unexpected_case = BashOperator( task_id='handle_unexpected_abnormal_count', bash_command='echo "当日异常记录数非预期,需人工检查"' ) # 设置分支依赖 branch_decider >> [success_report, failed_report, unexpected_case]
说明:
- 把
your_mysql_connection_id替换为Airflow中已配置的数据库连接ID BashOperator可根据实际业务替换为PythonOperator或其他自定义算子
2. 自定义Python脚本(配合Crond调度)
如果不需要复杂的调度平台,直接写Python脚本+系统定时任务即可:
import mysql.connector from datetime import datetime def run_success_report(): # 实现sending_reports_if_success的业务逻辑 print("执行正常状态报告发送任务") def run_failed_report(): # 实现sending_reports_if_failed的业务逻辑 print("执行异常状态报告发送任务") def main(): # 获取任务启动日期,按需调整格式 task_date = datetime.now().strftime('%Y-%m-%d') # 数据库连接配置 db_config = { 'host': 'your_db_host', 'user': 'your_db_user', 'password': 'your_db_password', 'database': 'your_db_name' } try: conn = mysql.connector.connect(**db_config) cursor = conn.cursor() query = """ SELECT COUNT(*) FROM r_summary rs WHERE rs.status = 'Not Ok' AND rs.created_date = %s; """ cursor.execute(query, (task_date,)) abnormal_count = cursor.fetchone()[0] if abnormal_count == 1: run_failed_report() elif abnormal_count == 0: run_success_report() else: print(f"警告:当日异常记录数为{abnormal_count},超出预期范围") finally: if conn.is_connected(): cursor.close() conn.close() if __name__ == "__main__": main()
然后用Crond设置每日调度,比如添加以下定时规则:
0 9 * * * /usr/bin/python3 /path/to/your/script.py >> /var/log/report_task.log 2>&1
3. 通用调度工具(如XXL-JOB、Jenkins)实现思路
不管用哪种调度平台,核心步骤一致:
- 新增前置判断任务:执行指定SQL查询
count值,并将结果存入上下文或变量 - 配置分支逻辑:
- 当
count=0时,触发sending_reports_if_success任务 - 当
count=1时,触发sending_reports_if_failed任务
- 当
- 补充异常分支:处理
count不为0或1的情况(比如多条异常记录),避免逻辑遗漏
内容的提问来源于stack exchange,提问作者Umid Umaraliev
相关产品推荐
相关产品推荐

