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
相关产品推荐
相关产品推荐

