Airflow任务失败时如何保存进度并实现断点续跑?
基于Airflow+Python的单任务断点续跑方案
方案1:优化XCom状态存储(优先推荐)
XCom是Airflow原生的任务内状态存储机制,和Task Instance(TI)强绑定,天然隔离不同DAG Run的状态。你之前遇到的触发表单DAG冲突问题,核心解决思路是明确关联当前TI的XCom数据,而非全局读取。
实现代码示例
from airflow.models import TaskInstance from airflow.utils.session import create_session from airflow.operators.python import PythonOperator from airflow import DAG from datetime import datetime def batch_task_resume(**context): ti: TaskInstance = context["ti"] dag_run = context["dag_run"] # 从触发表单的conf中获取自定义参数(如果有) task_params = dag_run.conf or {} # 定义25个待执行步骤(可替换为从conf读取) steps = [f"run_command_{i}" for i in range(1, 26)] # 读取当前TI已完成的步骤,默认空列表 completed_steps = ti.xcom_pull(key="completed_steps", task_ids=ti.task_id) or [] for step in steps: if step in completed_steps: print(f"跳过已完成步骤: {step}") continue # 执行当前步骤(替换为实际命令执行逻辑,如subprocess.run) print(f"执行步骤: {step}") # 模拟命令执行成功(实际需捕获异常,失败则抛出中断) # if execute_command(step) != 0: # raise Exception(f"步骤{step}执行失败") # 标记步骤完成并立即写入XCom completed_steps.append(step) ti.xcom_push(key="completed_steps", value=completed_steps) # 强制提交会话,确保状态实时写入Airflow元数据库 with create_session() as session: session.merge(ti) session.commit() # 定义带UI触发表单的DAG with DAG( dag_id="batch_resume_dag", start_date=datetime(2024, 1, 1), schedule_interval=None, # 手动触发 params={ # 定义UI触发时的表单字段 "batch_name": {"type": "string", "label": "批次名称", "default": "default_batch"} } ) as dag: batch_task = PythonOperator( task_id="batch_task", python_callable=batch_task_resume, provide_context=True, execution_timeout=None # 根据实际任务时长调整 )
关键注意点
- 每次步骤执行成功后立即写入XCom并提交会话,避免中途失败导致状态丢失
- 重跑任务时,需确保不清除该TI的XCom数据(Airflow默认重跑会保留XCom,除非手动清除)
- 若需区分不同触发批次,可结合
dag_run.conf中的唯一标识(如batch_id)作为XCom键的一部分,进一步隔离状态
方案2:利用Airflow Variable做隔离式状态存储
Variable是Airflow的全局键值存储,通过给每个DAG Run生成唯一前缀,可实现状态隔离,适合对XCom有特殊限制的场景。
实现代码示例
from airflow.models import Variable, TaskInstance from airflow.utils.session import create_session from airflow.operators.python import PythonOperator from airflow import DAG from datetime import datetime def batch_task_resume(**context): ti: TaskInstance = context["ti"] dag_run = context["dag_run"] run_id = dag_run.run_id # 生成唯一状态前缀,避免不同DAG Run冲突 state_prefix = f"batch_task_{run_id}_" steps = [f"run_command_{i}" for i in range(1, 26)] for step in steps: state_key = f"{state_prefix}{step}" # 检查步骤是否已完成 if Variable.get(state_key, default_var=None) == "completed": print(f"跳过已完成步骤: {step}") continue # 执行步骤逻辑 print(f"执行步骤: {step}") # 模拟执行成功 # 标记步骤完成 Variable.set(state_key, "completed") with create_session() as session: session.commit() # 任务完成后清理当前Run的状态变量(可选,避免冗余存储) for step in steps: state_key = f"{state_prefix}{step}" Variable.delete(state_key) with DAG( dag_id="batch_resume_var_dag", start_date=datetime(2024, 1, 1), schedule_interval=None, params={"batch_name": {"type": "string", "label": "批次名称"}} ) as dag: batch_task = PythonOperator( task_id="batch_task", python_callable=batch_task_resume, provide_context=True )
关键注意点
- 必须用
run_id或触发表单中的唯一标识作为前缀,防止全局Variable冲突 - 任务完成后建议清理状态变量,避免Airflow元数据库存储膨胀
- 确保当前Airflow角色有Variable的读写权限
方案3:基于Task Instance自定义备注存储状态(轻量备选)
利用Task Instance的note字段存储已完成步骤的序列化数据(如逗号分隔字符串),适合极简场景:
def batch_task_resume(**context): ti: TaskInstance = context["ti"] steps = [f"run_command_{i}" for i in range(1, 26)] # 解析note中的已完成步骤 completed_steps = ti.note.split(",") if ti.note else [] for step in steps: if step in completed_steps: print(f"跳过已完成步骤: {step}") continue print(f"执行步骤: {step}") # 模拟执行成功 completed_steps.append(step) ti.note = ",".join(completed_steps) with create_session() as session: session.merge(ti) session.commit()
局限性
note字段长度有限,不适合存储大量步骤的状态- 可读性较差,仅适合步骤数量少的场景
内容的提问来源于stack exchange,提问作者HFX
相关产品推荐
相关产品推荐

