如何在Airflow的on_failure_callback中访问任务参数实现回滚
在Airflow失败回调中获取任务参数的方法
针对你遇到的问题,有两种可靠方式可以在on_failure_callback的rollback函数中获取传入任务的参数:
方法一:通过TaskInstance直接读取任务参数
Airflow的context字典中包含ti(TaskInstance)对象,通过它可以访问任务实例的task属性,进而获取位置参数或关键字参数:
位置参数场景(如你的示例)
修改rollback函数如下:
def rollback(context: dict): ti = context["ti"] # 从op_args中提取第一个位置参数 task_argument = ti.task.op_args[0] print(f"获取到的task_argument值:{task_argument}")
关键字参数场景
如果调用任务时使用关键字参数(如example_task(task_argument="参数值")),则调整为:
def rollback(context: dict): ti = context["ti"] # 从op_kwargs中提取指定关键字的参数 task_argument = ti.task.op_kwargs["task_argument"] print(f"获取到的task_argument值:{task_argument}")
方法二:通过XCom获取自动推送的任务参数
Airflow会自动将@task装饰的任务参数推送到XCom中:
- 位置参数存储在
task_args键下 - 关键字参数存储在
task_kwargs键下
位置参数场景
def rollback(context: dict): ti = context["ti"] # 拉取XCom中的位置参数列表 task_args = ti.xcom_pull(key="task_args", task_ids=ti.task_id) task_argument = task_args[0] print(f"获取到的task_argument值:{task_argument}")
关键字参数场景
def rollback(context: dict): ti = context["ti"] # 拉取XCom中的关键字参数字典 task_kwargs = ti.xcom_pull(key="task_kwargs", task_ids=ti.task_id) task_argument = task_kwargs["task_argument"] print(f"获取到的task_argument值:{task_argument}")
注意事项
- 以上方法适用于Airflow 2.x版本,
@task装饰器是该版本引入的特性。 - 若任务参数为复杂自定义对象,需确保其可序列化(符合Airflow XCom的序列化规则),否则可能无法正确读取。
内容的提问来源于stack exchange,提问作者Izaak Cornelis
相关产品推荐
相关产品推荐

