You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 04:53:18