You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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)实现思路

不管用哪种调度平台,核心步骤一致:

  1. 新增前置判断任务:执行指定SQL查询count值,并将结果存入上下文或变量
  2. 配置分支逻辑:
    • 当count=0时,触发sending_reports_if_success任务
    • 当count=1时,触发sending_reports_if_failed任务
  3. 补充异常分支:处理count不为0或1的情况(比如多条异常记录),避免逻辑遗漏

内容的提问来源于stack exchange,提问作者Umid Umaraliev

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 05:59:53