Airflow 2.4.3周日清理XCom的DAG报语法错误求助
问题:Airflow DAG启动时出现十进制整数前导零语法错误
在Airflow 2.4.3版本中创建每周日运行的XCom清理DAG,使用Cron表达式5 16 * * 0调度,运行时抛出语法错误,提示十进制整数字面量不允许前导零,需用0o前缀表示八进制整数,错误定位在start_date参数的datetime(2023,12,03,16,05)处。
错误栈
Broken DAG: [/usr/local/airflow/dags/xcom_clear_30_days.py] Traceback (most recent call last): File "<frozen importlib._bootstrap_external>", line 947, in source_to_code File "<frozen importlib._bootstrap>", line 241, in _call_with_frames_removed File "/usr/local/airflow/dags/xcom_clear_30_days.py", line 10 with DAG(dag_id="cleanup_30_days_older_xcom_values", schedule_interval="5 16 * * 0", start_date=datetime(2023,12,03, 16 ,05),tags=["cleanup_30_days"]) as dag: ^ SyntaxError: leading zeros in decimal integer literals are not permitted; use an 0o prefix for octal integers
原代码
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_30_days_older_xcom_values", schedule_interval="5 16 * * 0", start_date=datetime(2023,12,03, 16 ,05),tags=["cleanup_30_days"]) 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 the time until the xcom cleared {ts_limit}") print(f"printing session {session}") ##session.query(XCom).filter(XCom.execution_date <= ts_limit).delete() session.query(XCom).filter(XCom.execution_date <= ts_limit).delete(synchronize_session='fetch') 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
解决方法
Python 3禁止十进制整数使用前导零(比如03、05),这是语法规范问题。只需将datetime参数中的前导零去掉即可:
- 把
datetime(2023,12,03,16,05)改为datetime(2023,12,3,16,5)
另外,Airflow官方推荐使用days_ago()来设置start_date,可以避免手动编写日期时的格式错误,比如改为start_date=days_ago(1)。
修正后的代码
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 # 方案1:修正datetime的前导零问题 with DAG(dag_id="cleanup_30_days_older_xcom_values", schedule_interval="5 16 * * 0", start_date=datetime(2023,12,3,16,5), tags=["cleanup_30_days"]) as dag: # 方案2:使用days_ago更简洁(推荐) # with DAG(dag_id="cleanup_30_days_older_xcom_values", schedule_interval="5 16 * * 0", start_date=days_ago(1), tags=["cleanup_30_days"]) as dag: # cleanup_xcom @provide_session def cleanup_xcom(session=None, **context): dag = context["dag"] dag_id = dag._dag_id # 删除30天前的XCom数据 ts_limit = datetime.now(timezone.utc) - timedelta(days=30) print(f"清理截止时间:{ts_limit}") print(f"数据库会话:{session}") session.query(XCom).filter(XCom.execution_date <= ts_limit).delete(synchronize_session='fetch') 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
内容的提问来源于stack exchange,提问作者dewdrops
相关产品推荐
相关产品推荐

