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
相关产品推荐
相关产品推荐

