Airflow任务与DAG自动重试配置求助:失败后重试两次再终止
Airflow任务与DAG重试配置方案(针对SFTP到S3大文件传输超时场景)
一、单个Task失败自动重试配置
针对你的77个传输Task,只需在Task定义时添加重试参数,即可实现失败后自动重试2次,全部失败后再标记Task失败:
核心参数说明
retries=2:设置重试次数为2次(加上首次执行,每个Task最多尝试3次)retry_delay:设置两次重试的间隔时间,针对AWS超时场景,建议设置5-10分钟的延迟,避免立刻重试再次触发超时- 可选:
exception_retry_rules:仅针对AWS超时相关异常重试,避免无关错误消耗资源
代码示例(以PythonOperator为例,其他Operator同理)
from airflow.operators.python import PythonOperator from datetime import timedelta def sftp_to_s3_transfer(): # 你的SFTP到S3传输逻辑 pass transfer_task = PythonOperator( task_id="sftp_to_s3_task_01", python_callable=sftp_to_s3_transfer, retries=2, retry_delay=timedelta(minutes=5), # 仅捕获AWS超时类异常重试 exception_retry_rules=lambda e: "TimeoutError" in str(e) or "AWS" in str(e).upper() )
如果是批量创建77个Task,可以把这些重试参数封装到default_args里批量应用:
task_default_args = { "retries": 2, "retry_delay": timedelta(minutes=5), "exception_retry_rules": lambda e: "TimeoutError" in str(e) or "AWS" in str(e).upper() } # 循环创建11*7个Task for i in range(11): for j in range(7): task = PythonOperator( task_id=f"transfer_task_{i}_{j}", python_callable=sftp_to_s3_transfer, default_args=task_default_args ) # 配置Task依赖(如果有)
二、DAG级失败自动重试配置
当所有Task重试后仍有失败导致DAG终止时,通过DAG级配置实现整个DAG自动重试2次:
核心参数说明
- 在DAG的
default_args中设置retries=2和retry_delay,对整个DAG生效 - 确保
catchup=False(如果不需要补历史任务),避免重试触发历史实例
- 在DAG的
代码示例
from airflow import DAG from datetime import datetime, timedelta dag_default_args = { "owner": "airflow", "start_date": datetime(2024, 1, 1), "retries": 2, # DAG级重试2次 "retry_delay": timedelta(minutes=10) # DAG重试间隔10分钟 } with DAG( dag_id="sftp_s3_batch_transfer_dag", default_args=dag_default_args, schedule_interval="@daily", catchup=False ) as dag: # 这里放置你的77个Task定义及依赖配置
注意事项
- Task级重试优先于DAG级重试:单个Task会先完成自己的2次重试,全部失败后才会触发DAG的重试逻辑
- 重试间隔需根据实际场景调整:大文件传输失败后,过短的间隔可能再次遇到AWS超时,建议根据AWS服务状态和网络情况设置合理延迟
- 可通过Airflow UI查看重试记录:在Task实例详情页的"Logs"或"Retry"标签下查看每次重试的执行情况
内容的提问来源于stack exchange,提问作者Jomy
相关产品推荐
相关产品推荐

