Airflow 2.3.4:如何在任务故意失败的回调中标记为成功?
解决Airflow 2.3.4中失败回调标记任务为成功的问题
你的代码无法生效的核心原因是:直接给task_instance.state赋值只是修改了内存中的对象属性,没有将状态变更同步到Airflow的元数据库,所以调度器和Web UI都不会识别这个状态变化。
要实现需求,需要使用Airflow提供的官方状态更新方法,并确保变更持久化到数据库,具体代码如下:
from airflow.exceptions import AirflowException from airflow.utils.state import State from airflow.operators.python import PythonOperator from airflow.settings import Session def intentional_failure(): raise AirflowException("this is a dummy failure") def handle_failure(context): ti = context['task_instance'] # 调用set_state方法更新状态,该方法会处理Airflow内部的状态变更逻辑 ti.set_state(State.SUCCESS) # 获取任务实例关联的数据库会话并提交,确保变更写入元数据库 session = Session.object_session(ti) if session: session.commit() dummy_failure = PythonOperator( task_id="intentional_failure", python_callable=intentional_failure, on_failure_callback=handle_failure )
关键说明
set_state()是Airflow官方提供的任务状态更新方法,它不仅会修改状态属性,还会同步更新任务的end_date、duration等关联字段,符合Airflow的状态管理逻辑。- 通过
Session.object_session(ti)获取任务实例绑定的数据库会话并提交,能确保状态变更被持久化到元数据库,避免仅在内存中修改状态的问题。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

