Airflow集成Snowflake:SQL文件最后查询仅返回单行而非全部行问题
问题:Snowflake批量执行SQL后仅返回单行结果,预期获取最后一条查询的所有行
问题背景
需求为获取Snowflake SQL文件中最后一条查询语句的所有行数据,但当前代码仅返回单行内容,不符合预期。
现有代码实现
def get_ids(self): snflk_hook = SnowflakeHook(snowflake_conn_id=self.snflk_conn_id) with closing(snflk_hook.get_conn()) as conn: with closing(conn.cursor(DictCursor)) as cur: cur.execute(self.sql) ids = [] for row in cur: ids.append(row) return mongo_ids
执行的SQL脚本
BEGIN CREATE OR REPLACE TEMP TABLE ABC_ID AS SELECT * FROM XYZ; SELECT ID FROM ABC_ID; END
问题分析与解决方法
问题原因
当执行包含多个语句的Snowflake脚本(如BEGIN/END包裹的代码块)时,cur.execute()会生成多个结果集。默认下游标只会遍历第一个结果集(此处为CREATE OR REPLACE TEMP TABLE的执行返回),因此只能拿到单行数据,而最后一条SELECT ID的结果集并未被处理。
修正后的代码
def get_ids(self): snflk_hook = SnowflakeHook(snowflake_conn_id=self.snflk_conn_id) with closing(snflk_hook.get_conn()) as conn: with closing(conn.cursor(DictCursor)) as cur: cur.execute(self.sql) ids = [] # 遍历所有结果集,直到处理完最后一个 while True: # 提取当前结果集的所有行 for row in cur: ids.append(row) # 切换到下一个结果集,无后续结果集则退出循环 if not cur.nextset(): break return ids
额外注意点
原代码存在笔误:定义的结果列表是ids,但返回时写的是mongo_ids,需修正为返回ids,否则会引发变量未定义的错误。
内容的提问来源于stack exchange,提问作者Kar
相关产品推荐
相关产品推荐

