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

Airflow回填任务随机出现「任务状态被外部修改」错误求助

Cloud Composer中Deferrable模式任务回填时的状态不一致问题

我在Google Cloud Composer环境中使用Apache Airflow运行DAG,由于任务运行时间较长,几乎所有任务都以deferrable模式执行。

在针对单个DAG运行实例执行回填任务时,会随机出现以下错误日志:

[2025-08-27T11:47:02.620+0000] {backfill_job_runner.py:323} ERROR - The executor reported that the task instance <TaskInstance: my_dag.my_task backfill__2025-07-15T03:00:00+02:00 [queued]> finished with state success, but the task instance's state attribute is queued.

该错误同样会在{taskinstance.py:3314}中上报。除非在回填命令中指定--disable-retry禁用重试,否则任务会自动重试。

在Airflow UI中,任务实例日志显示任务已成功完成:

[2025-08-27, 13:46:59 CEST] {taskinstance.py:352} INFO - Marking task as SUCCESS. dag_id=my_dag, task_id=my_task, run_id=backfill__2025-07-15T03:00:00+02:00, execution_date=20250715T010000, start_date=20250827T112744, end_date=20250827T114659
[2025-08-27, 13:46:59 CEST] {local_task_job_runner.py:266} INFO - Task exited with return code 0
[2025-08-27, 13:46:59 CEST] {taskinstance.py:3903} INFO - 1 downstream tasks scheduled from follow-on schedule check

但UI中任务本身显示为失败状态。

执行的回填命令

airflow dags backfill \
  --conf='{"test":true}' \
  --start-date=2025-07-15T03:00:00+0200 \
  --end-date=2025-07-15T04:00:00+0200 \
  --reset-dagruns \
  --yes \
  my_dag

DAG简化代码

from datetime import datetime
from airflow.timetables.trigger import CronTriggerTimetable
from airflow.decorators import dag
from airflow.models.baseoperator import chain

from airflow.providers.google.cloud.operators.vertex_ai.batch_prediction_job import CreateBatchPredictionJobOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator

@dag(
    start_date=datetime(2025, 6, 1),
    catchup=False,
    schedule=CronTriggerTimetable("00 03 15 * *", timezone="Europe/Berlin"),
    template_searchpath=["/path/to/templates/"],
    params={"test": False},
)
def my_data_pipeline():

    bq_task_initial = BigQueryInsertJobOperator(
        task_id="create_table_initial_data",
        configuration={
            "query": {
                "query": "{% include 'data/sql/query_initial_data.sql' %}",
            }
        },
        deferrable=True,
        poll_interval=30
    )

    ml_task = CreateBatchPredictionJobOperator(
        task_id="run_batch_prediction_job",
        batch_prediction_job=prediction_job,
        deferrable=True,
        poll_interval=30,
        # 更多初始化字段
        # ..
    )

    bq_task_final = BigQueryInsertJobOperator(
        task_id="create_table_final_processing",
        configuration={
            "query": {
                "query": "{% include 'data/sql/query_final.sql' %}",
            }
        },
        deferrable=True,
        poll_interval=30
    )

    chain(
        bq_task_initial,
        ml_task,
        bq_task_final
    )

my_data_pipeline()

Cloud Composer环境配置

Composer Version:   2.11.1
Airflow Version:    2.10.2

Workloads configuration
  Scheduler:        1 scheduler with 0.5 vCPU, 2 GB memory, 1 GB storage
  Triggerer:        1 triggerer with 0.5 vCPU, 0.5 GB memory, 1 GB storage
  Web server:       0.5 vCPU, 2 GB memory, 1 GB storage
  Worker:           Autoscaling between 1 and 3 workers, with 0.5 vCPU, 2 GB memory, 1 GB storage each

Core infrastructure
  Environment size: Small

Airflow核心配置

core            dagbag_import_timeout           120
core            dag_file_processor_timeout      300
scheduler       task_queued_timeout             2400
celery          worker_concurrency              3

已尝试的解决方法

  • 调整上述Airflow配置参数,无效果
  • 回填时添加--donot-pickle参数,无效果

疑问

该问题是否特定于回填任务和/或deferrable模式任务?是否有其他人遇到类似问题?可能的解决方案是什么?


内容的提问来源于stack exchange,提问作者Murphy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 11:27:26