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

Airflow升级后SnowflakeCheckOperator报错:'dict'无'sfqid'属性

问题分析

升级MWAA(Airflow 2.2.2 → 2.4.3)后,使用SnowflakeCheckOperator或dbapi_hook.get_first()时触发AttributeError: 'dict' object has no attribute 'sfqid',但SnowflakeOperator或直接调用dbapi_hook.run()无异常。

从报错栈可知,问题出在airflow/providers/snowflake/hooks/snowflake.py的run方法中:当传入handler参数时(比如get_first()会传入fetch_one_handler),方法会先执行handler(cur)返回处理后的结果(这里是查询结果的字典),但后续代码仍尝试访问cur.sfqid,此时cur已经被替换成handler的返回值(字典),而非Snowflake游标对象,导致报错。

而SnowflakeOperator调用run()时未传入handler,cur保持为游标对象,因此能正常访问sfqid。

解决方案

方案1:降级Snowflake Provider版本

将apache-airflow-providers-snowflake从3.0.0降级到2.3.0,该版本的run方法逻辑未引入此问题。修改requirements.txt:

apache-airflow-providers-snowflake==2.3.0

方案2:自定义修复后的SnowflakeHook

编写自定义Hook,调整run方法逻辑,确保在调用handler前先获取sfqid:

from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook

class FixedSnowflakeHook(SnowflakeHook):
    def run(
        self,
        sql: str | list[str],
        parameters: dict | list | None = None,
        autocommit: bool = False,
        handler: callable | None = None,
        split_statements: bool = False,
        return_last: bool = True,
    ):
        self.log.info("Running statement: %s, parameters: %s", sql, parameters)
        with self.get_conn() as conn:
            if split_statements:
                statements = self.split_sql(sql)
            else:
                statements = [sql]

            if autocommit:
                conn.autocommit = autocommit

            query_ids = []
            results = []
            for stmt in statements:
                with conn.cursor() as cur:
                    self._run_command(cur, stmt, parameters)
                    # 先获取sfqid再调用handler
                    query_id = cur.sfqid
                    query_ids.append(query_id)
                    if handler is not None:
                        result = handler(cur)
                        results.append(result)
                    self.log.info("Statement execution info - %s", query_id)

            if handler is not None:
                return results[-1] if return_last else results
            return query_ids[-1] if return_last else query_ids

然后在DAG中使用自定义Hook:

from airflow.providers.common.sql.operators.sql import SQLCheckOperator

check_snowflake_table = SQLCheckOperator(
    task_id="check_snowflake_table",
    sql="SELECT COUNT(*) FROM DEV_MONITORING.NZO_STATE.TBL_DAG_MESSAGE_QUEUE WHERE DAG_NAME = 'TEST_DAG'",
    hook=FixedSnowflakeHook(snowflake_conn_id="snowflake_conn-non-live"),
)

方案3:替代SnowflakeCheckOperator

改用SnowflakeOperator执行查询,手动判断结果是否符合检查条件:

from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from airflow.operators.python import PythonOperator

def check_snowflake_result(**context):
    ti = context["ti"]
    query_result = ti.xcom_pull(task_ids="run_check_query")
    # 根据需求调整判断逻辑,示例为检查COUNT(*)是否大于0
    if int(query_result[0][0]) == 0:
        raise ValueError("DAG名称未在表中找到")

with DAG(...) as dag:
    run_check_query = SnowflakeOperator(
        task_id="run_check_query",
        snowflake_conn_id="snowflake_conn-non-live",
        sql="SELECT COUNT(*) FROM DEV_MONITORING.NZO_STATE.TBL_DAG_MESSAGE_QUEUE WHERE DAG_NAME = 'TEST_DAG'",
        do_xcom_push=True,
    )

    validate_result = PythonOperator(
        task_id="validate_result",
        python_callable=check_snowflake_result,
        provide_context=True,
    )

    run_check_query >> validate_result

内容的提问来源于stack exchange,提问作者CallumO

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 16:53:17