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

Airflow中使用execute_string执行Snowflake多SQL语句报错排查

问题分析与修复方案

错误原因

你遇到的UnboundLocalError: local variable 'sf' referenced before assignment及潜在的方法调用错误,由两个核心问题导致:

  1. execute_string未成为SnowflakeHook类的成员方法:原代码中该方法定义在类外部,不属于SnowflakeHook的实例方法,即使sf初始化成功,调用时也会触发AttributeError。
  2. Hook初始化未处理异常:如果SnowflakeHook的__init__方法中(如Snowflake连接失败、凭证错误)抛出异常,sf = SnowflakeHook(sf_conn_id)会中断执行,sf变量未被赋值,后续调用sf.execute_string时就会触发未绑定变量错误。

修复步骤

1. 修正SnowflakeHook类的execute_string方法

将execute_string缩进至类内部,成为类的成员方法,同时优化多语句结果的获取逻辑(默认fetchall()仅返回最后一条语句结果):

from airflow.hooks.base_hook import BaseHook
class SnowflakeHook(BaseHook):
    def __init__(self, sf_conn_id, warehouse=None, source=None):
        """Snowflake hook init.
        Arguments:
        sf_conn_id: Airflow connection ID providing Snowflake credential.
        warehouse: Warehouse override.  Defaults to "warehouse" key in the connection extra section.
        source: Hook source to pass to the parent class constructor.
        """
        import snowflake.connector
        super().__init__(source)
        tmp_conn = self.get_connection(sf_conn_id)
        self.user = tmp_conn.login
        self.password = tmp_conn.password
        extras = tmp_conn.extra_dejson
        self.account = extras.get("account", None)
        self.warehouse = warehouse if warehouse is not None else extras.get("warehouse", None)
        self.database = extras.get("database", None)
        self.schema = tmp_conn.schema
        self.role=extras.get("role", None)
        self.conn = snowflake.connector.connect(
            user=self.user,
            password=self.password,
            account=self.account,
            warehouse=self.warehouse,
            database=self.database,
            schema=self.schema,
            role=self.role,
        )
    
    def execute(self, qry):
        """Execute a sql statement"""
        cs = self.conn.cursor()
        try:
            cs.execute(qry)
            return cs.fetchall()
        finally:
            cs.close()
    
    # 修正:缩进至类内部,成为实例方法
    def execute_string(self, sql_str):
        """Executes one or more SQL statements separated by semi-colons"""
        cs = self.conn.cursor()
        try:
            # 遍历所有语句的结果,收集全部返回值
            all_results = []
            for result in cs.execute_string(sql_str):
                all_results.extend(result.fetchall())
            return all_results
        finally:
            cs.close()

2. 优化DAG调用函数的异常处理

添加异常捕获和资源清理逻辑,确保变量正确赋值、连接及时关闭:

def delete(sf_conn_id):
    from custom.hooks.snowflake_hook import SnowflakeHook
    table_list_sql = """
                    create temporary table dim_expired  
                        as select * from expired_insert_stream; 
                    select distinct id  
                      from dim_expired
                    """
    sf = None
    try:
        sf = SnowflakeHook(sf_conn_id)
        print(table_list_sql)
        res = sf.execute_string(table_list_sql)
        print("查询结果:", res)
    except Exception as e:
        print(f"执行出错: {str(e)}")
        raise  # 重新抛出异常,让Airflow捕获并标记任务失败
    finally:
        # 确保Snowflake连接关闭,避免资源泄漏
        if sf and hasattr(sf, 'conn'):
            sf.conn.close()

额外建议

  • 优先使用Airflow官方提供的SnowflakeHook,官方实现已支持多语句执行、连接池、异常处理等成熟功能,无需自定义Hook。
  • 若必须自定义,建议添加连接关闭的专用方法,或使用上下文管理器(__enter__/__exit__)来管理连接生命周期。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 17:31:00