Airflow 2.4.3清理30天前XCom数据报错:Xcom/execution_date未定义
清理Airflow XCom时的"Xcom is not defined"及"execution_date not defined"问题解决
问题背景
尝试清理execution_date早于30天的XCom值,使用Airflow 2.4.3版本,运行DAG时先后出现两个报错:
- 提示
Xcom is not defined - 移除
Xcom前缀后又提示execution_date not defined
原代码
from airflow.models import DAG from airflow.utils.db import provide_session from airflow.models import XCom from airflow.operators.python import PythonOperator from airflow.operators.dummy import DummyOperator from airflow.utils.dates import days_ago from datetime import datetime, timedelta,timezone from sqlalchemy import func with DAG(dag_id="cleanup_xcom_demo", schedule_interval=None, start_date=days_ago(2)) as dag: # cleanup_xcom @provide_session def cleanup_xcom(session=None, **context): dag = context["dag"] dag_id = dag._dag_id # It will delete all xcom of the dag_id ts_limit = datetime.now(timezone.utc) - timedelta(days=30) print(f"printing session {session}") session.query(XCom).filter(Xcom.execution_date <= ts_limit ).delete() clean_xcom = PythonOperator( task_id="clean_xcom", python_callable = cleanup_xcom, provide_context=True, # dag=dag ) start = DummyOperator(task_id="start") end = DummyOperator(task_id="end", trigger_rule="none_failed") start >> clean_xcom >> end
报错信息
Xcom is not defined when i remove xcom from xcom.execution_date and keep only execution_date in filter then it throws error as execution_date not defined
问题分析与解决
1. 大小写错误导致Xcom is not defined
导入的模型是XCom(大写C),但代码过滤条件里写的是Xcom.execution_date(小写c),Python是大小写敏感语言,因此触发未定义报错。
2. 缺少模型关联导致execution_date not defined
移除Xcom.前缀后,execution_date没有关联到XCom模型,SQLAlchemy无法识别这个字段,因此报错。
额外优化:限定当前DAG的XCom
原代码未过滤dag_id,会删除所有DAG中符合时间条件的XCom,建议加上dag_id过滤,只清理当前DAG的XCom数据。
修正后的代码
from airflow.models import DAG from airflow.utils.db import provide_session from airflow.models import XCom from airflow.operators.python import PythonOperator from airflow.operators.dummy import DummyOperator from airflow.utils.dates import days_ago from datetime import datetime, timedelta, timezone with DAG(dag_id="cleanup_xcom_demo", schedule_interval=None, start_date=days_ago(2)) as dag: @provide_session def cleanup_xcom(session=None, **context): dag_id = context["dag"].dag_id # 直接使用公开属性dag_id,避免访问私有变量 ts_limit = datetime.now(timezone.utc) - timedelta(days=30) # 过滤当前DAG且execution_date早于30天的XCom session.query(XCom).filter( XCom.dag_id == dag_id, XCom.execution_date <= ts_limit ).delete() clean_xcom = PythonOperator( task_id="clean_xcom", python_callable=cleanup_xcom, provide_context=True, ) start = DummyOperator(task_id="start") end = DummyOperator(task_id="end", trigger_rule="none_failed") start >> clean_xcom >> end
补充说明
- Airflow 2.x中
provide_context=True在PythonOperator中可被op_kwargs替代,但2.4.3版本仍兼容该参数。 - 操作数据库前建议在测试环境验证逻辑,避免误删数据。
内容的提问来源于stack exchange,提问作者dewdrops
相关产品推荐
相关产品推荐

