Airflow中session.query()获取任务最后执行时间时生成SQL语法错误
Airflow 2.4.2查询任务最近成功执行时间的SQL语法错误问题解决
问题场景
在容器部署的Airflow 2.4.2环境中,开发任务时需要获取指定DAG任务的最近N次成功执行时间,参考方案编写了基于SQLAlchemy查询PostgreSQL数据库的函数,但运行时触发SQL语法错误,错误显示ORDER BY子句的子查询内DESC被错误放置在WHERE条件末尾。
报错信息
核心报错:
sqlalchemy.exc.ProgrammingError: (psycopg2.errors.SyntaxError) syntax error at or near "DESC" LINE 5: ..._id = task_instance.run_id AND dag_run.execution_date DESC) ^
生成的错误SQL片段:
ORDER BY EXISTS ( SELECT 1 FROM dag_run WHERE dag_run.dag_id = task_instance.dag_id AND dag_run.run_id = task_instance.run_id AND dag_run.execution_date DESC )
问题根源
这个问题是Airflow 2.4.x版本的已知Bug,并非SQLAlchemy的问题。Airflow的TaskInstance模型中,execution_date字段关联了dag_run表的对应字段,当通过SQLAlchemy查询并排序TaskInstance.execution_date时,Airflow的模型关联逻辑会自动生成关联dag_run的子查询,但在2.4.x版本中,排序关键字DESC被错误地嵌入到了子查询的WHERE条件中,导致SQL语法错误。该Bug在Airflow 2.5.0及后续版本已被修复。
解决方法
方案1:升级Airflow版本
直接将Airflow升级到2.5.0或更高版本,这是最彻底的解决方式,Airflow官方在该版本修复了这个模型关联排序的Bug。
方案2:调整查询逻辑(无需升级)
如果无法升级Airflow,可以修改查询代码,避免触发错误的关联逻辑,以下是两种可行的修改方式:
方式A:直接查询execution_date字段
不加载完整的TaskInstance对象,仅查询需要的execution_date字段,避免触发关联dag_run的自动逻辑:
from typing import List, Optional from airflow.utils.state import State from airflow.settings import Session from airflow.models.taskinstance import TaskInstance def last_execution_date( dag_id: str, task_id: str, n: int, session: Optional[Session] = None ) -> List[str]: session = session or Session() query = ( session.query(TaskInstance.execution_date) .filter( TaskInstance.dag_id == dag_id, TaskInstance.task_id == task_id, TaskInstance.state == State.SUCCESS ) .order_by(TaskInstance.execution_date.desc()) .limit(n) ) # 将datetime对象转换为字符串格式返回 execution_dates = [row[0].isoformat() for row in query.all()] return execution_dates
方式B:使用原生SQL查询
绕过SQLAlchemy的模型关联逻辑,直接执行原生SQL语句:
from typing import List, Optional from airflow.settings import Session def last_execution_date( dag_id: str, task_id: str, n: int, session: Optional[Session] = None ) -> List[str]: session = session or Session() sql = """ SELECT execution_date FROM task_instance WHERE dag_id = %s AND task_id = %s AND state = 'success' ORDER BY execution_date DESC LIMIT %s """ result = session.execute(sql, (dag_id, task_id, n)) execution_dates = [row[0].isoformat() for row in result.fetchall()] return execution_dates
内容的提问来源于stack exchange,提问作者Mahmoud Saad
相关产品推荐
相关产品推荐

