如何在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
相关产品推荐
相关产品推荐

