You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 12:17:31