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

如何在Airflow中通过装饰器实现任务启停分支逻辑?

实现Airflow DAG的分支逻辑(基于errors字典长度)

核心实现方案

利用Airflow的@task.branch装饰器定义分支判断任务,通过检查errors字典的长度决定后续执行路径:

  • 当len(errors) == 0时,执行数据库插入任务
  • 当len(errors) > 0时,执行Telegram错误通知任务

完整代码示例

假设你已有API数据获取逻辑,以下是补充分支、数据库操作、告警任务的完整实现:

from airflow import DAG
from airflow.decorators import task
from airflow.operators.empty import EmptyOperator
from datetime import datetime
# 按需导入数据库连接、Telegram通知相关库

default_args = {
    'owner': 'your_team',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

with DAG(
    dag_id='api_sync_to_db_with_alert',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
) as dag:

    @task
    def fetch_api_data():
        # 替换为你的实际API数据获取逻辑
        # 示例返回结构:包含数据和错误字典
        data = {"user_id": 123, "user_name": "test"}
        errors = {}  # 模拟无错误场景;可改为 {"api_error": "超时"} 测试错误分支
        return {"data": data, "errors": errors}

    @task.branch
    def judge_execution_path(ti):
        # 从上游任务拉取结果
        api_result = ti.xcom_pull(task_ids='fetch_api_data')
        errors = api_result.get('errors', {})
        
        # 根据错误字典长度返回对应任务ID
        if len(errors) == 0:
            return 'insert_data_to_db'
        else:
            return 'send_telegram_error_alert'

    @task
    def insert_data_to_db(ti):
        api_result = ti.xcom_pull(task_ids='fetch_api_data')
        raw_data = api_result.get('data')
        # 此处保留数据库插入逻辑(无需添加数据转换,按需求后续补充)
        print(f"成功插入数据:{raw_data}")

    @task
    def send_telegram_error_alert(ti):
        api_result = ti.xcom_pull(task_ids='fetch_api_data')
        errors = api_result.get('errors')
        # 格式化错误文本
        error_content = "\n".join([f"{err_type}: {msg}" for err_type, msg in errors.items()])
        # 替换为你的Telegram机器人发送逻辑
        print(f"发送Telegram告警:{error_content}")

    # 定义任务依赖关系
    start = EmptyOperator(task_id='start')
    end = EmptyOperator(task_id='end', trigger_rule='none_failed_min_one_success')

    start >> fetch_api_data() >> judge_execution_path()
    judge_execution_path() >> [insert_data_to_db(), send_telegram_error_alert()] >> end

关键细节说明

  • @task.branch装饰的任务必须返回任务ID字符串或任务ID列表,Airflow会根据返回值触发对应分支
  • 通过ti.xcom_pull()完成任务间的数据传递,获取上游返回的data和errors
  • 最终的end任务设置trigger_rule='none_failed_min_one_success',确保无论哪个分支执行成功,DAG都能正常收尾
  • 可直接替换示例中的API获取、数据库操作、Telegram发送逻辑为你的实际业务代码

内容的提问来源于stack exchange,提问作者John Doe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 02:25:00