能否为Airflow单个TaskFlow @task设置重试机制?
Airflow 2.5.1 TaskFlow 模式下为特定任务配置重试的方法
你遇到的问题是因为把重试参数写在了任务函数的形参里,而非@task装饰器中,Airflow不会识别函数参数里的retries配置。当然可以为单个TaskFlow任务单独设置重试,无需修改全局默认参数,正确的写法如下:
基础重试配置
直接在@task装饰器中指定retries参数即可:
@task(retries=2) def test_retries(): raise ValueError("I failed, please retry") test_retries()
进阶重试配置
如果需要更精细的重试控制,比如重试间隔、指数退避策略等,也可以在装饰器中一并配置:
from datetime import timedelta @task( retries=2, retry_delay=timedelta(seconds=10), retry_exponential_backoff=True, max_retry_delay=timedelta(minutes=5) ) def test_retries(): raise ValueError("I failed, please retry") test_retries()
原代码不生效的原因
你之前把retries=2作为任务函数的参数,这只是普通的函数输入变量,Airflow的TaskFlow框架只会读取@task装饰器中传递的参数,将其映射到底层的PythonOperator对应的属性上,所以函数形参里的配置不会生效。
内容的提问来源于stack exchange,提问作者Noumenon
相关产品推荐
相关产品推荐

