Airflow 2.2.5执行DAG时出现TaskInstanceKey KeyError求助
Airflow 2.2.5
airflow dags test 执行时KeyError问题处理方案 问题背景
- 运行环境:Python 3.6 + Airflow 2.2.5
- 操作步骤:执行命令
airflow dags test airflow_report1_email 2022-08-30测试DAG - 错误现象:抛出KeyError,错误日志如下:
File "/home/airflow/.local/lib/python3.6/site-packages/airflow/utils/session.py", line 70, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/jobs/backfill_job.py", line 826, in _execute session=session, File "/home/airflow/.local/lib/python3.6/site-packages/airflow/utils/session.py", line 67, in wrapper return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/jobs/backfill_job.py", line 739, in _execute_dagruns session=session, File "/home/airflow/.local/lib/python3.6/site-packages/airflow/utils/session.py", line 67, in wrapper return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/jobs/backfill_job.py", line 634, in _process_backfill_task_instances self._update_counters(ti_status=ti_status) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/utils/session.py", line 70, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/jobs/backfill_job.py", line 216, in _update_counters ti_status.running.pop(reduced_key) KeyError: TaskInstanceKey(dag_id='airflow_report1_email', task_id='list_all_files', run_id='backfill__2022-08-30T00:00:00+00:00', try_number=12)
涉及DAG代码
from datetime import datetime, timedelta from textwrap import dedent from airflow import DAG from airflow.operators.bash import BashOperator with DAG( 'airflow_report1_email', schedule_interval='0 12 * * *', default_args={ 'depends_on_past': True, 'email': ['REDACTED'], 'email_on_failure': True, 'email_on_retry': True, 'retries': 1, 'retry_delay': timedelta(minutes=5), }, description='DAG to report the health of data', start_date=datetime(2022, 8, 30), catchup=False, tags=['example'], ) as dag: t1 = BashOperator( task_id='list_all_files', bash_command='/list_all_files.sh ', retries=1, ) t1.doc_md = dedent( """\ #### list_all_files This task downloads the (1) list of all files, (2) for each file for the past week checks if it's new, (3) downloads the new file. Estimated time to run: ~1h """ ) t2 = BashOperator( task_id='report_email', bash_command='/report1_email.sh ', retries=1, ) t2.doc_md = dedent( """\ #### report1_email This task computes the health of the newly downloaded data (last day of data). Estimated time to run: 5min """ ) dag.doc_md = """ This is a DAG to report the health of data (project) """ t1 >> t2
前提说明
已配置airflow用户权限,确认无操作权限问题;该问题为Airflow官方已知bug,在2.3.3及以上版本已修复。
解决方案
方案1:升级Airflow版本
这是最彻底的修复方式,将Airflow升级至2.3.3及以上版本即可解决该backfill任务计数器更新时的KeyError问题。注意Python 3.6仅兼容Airflow 2.5及以下版本(Airflow 2.6+要求Python 3.7+),升级时选择2.3.3~2.5.x区间内的稳定版本。
方案2:临时规避措施(不升级版本)
若暂时无法升级,可尝试以下操作:
- 清理目标DAG的历史任务实例:通过Airflow UI或命令行删除对应
run_id的TaskInstance记录,避免重试次数过多触发异常 - 临时调整重试配置:将任务的
retries设置为0,减少重试逻辑触发的计数器问题 - 改用单任务测试命令:使用
airflow tasks test airflow_report1_email list_all_files 2022-08-30替代airflow dags test,避免触发backfill流程中的计数器逻辑
内容的提问来源于stack exchange,提问作者martin
相关产品推荐
相关产品推荐

