Snowflake异步查询执行问题:SQL执行已取消
如何用Snowflake Python Connector异步执行查询并通过查询ID监控状态
我完全懂你处理海量数据时,想要异步提交查询、靠查询ID追踪状态的需求——这在Snowflake里是非常常见的实践,下面我给你整理完整的实现方案,直接就能用:
先准备好环境
首先确保你的Snowflake Python Connector是最新版本,避免兼容性问题:
pip install snowflake-connector-python --upgrade
核心实现思路
- 用
execute_async()方法提交异步查询:这个方法是非阻塞的,提交后立刻返回查询ID,查询会在Snowflake后台默默跑,不用一直占着连接 - 两种方式监控状态:要么用连接器自带的
get_query_status()快速查状态,要么查Snowflake系统视图QUERY_HISTORY获取更详细的执行信息 - 连接可以随时关闭:异步查询提交后,本地连接关了也不影响后台执行,后续查状态再重新连就行
完整代码示例
把你的代码片段补全后,完整的实现是这样的:
from __future__ import print_function import io, os, sys, time, datetime import snowflake.connector # 替换成你的Snowflake连接信息 SNOWFLAKE_CONFIG = { 'account': '你的账号(比如xxx.us-west-2)', 'user': '你的用户名', 'password': '你的密码', 'warehouse': '处理海量数据的仓库', 'database': '目标数据库', 'schema': '目标模式', 'role': '有权限的角色' } def submit_large_query(query_text): """提交异步查询,返回查询ID""" conn = None try: conn = snowflake.connector.connect(**SNOWFLAKE_CONFIG) cursor = conn.cursor() # 关键:用execute_async而非execute,直接拿到query_id query_id = cursor.execute_async(query_text) print(f"异步查询已提交!查询ID:{query_id}") return query_id finally: # 提交后就可以关连接,不影响后台执行 if conn: conn.close() def track_query_status(query_id): """根据查询ID检查状态,返回当前状态""" conn = None try: conn = snowflake.connector.connect(**SNOWFLAKE_CONFIG) cursor = conn.cursor() # 方式1:用连接器内置方法快速查状态 current_status = cursor.get_query_status(query_id) print(f"当前查询状态:{current_status}") # 方式2:查系统视图获取更详细信息(比如执行时长、扫描数据量) cursor.execute(""" SELECT QUERY_ID, QUERY_TEXT, STATUS, START_TIME, END_TIME, TOTAL_ELAPSED_TIME, SCANNED_ROWS FROM TABLE(INFORMATION_SCHEMA.QUERY_HISTORY()) WHERE QUERY_ID = %s """, (query_id,)) detail = cursor.fetchone() if detail: print(f"查询详情:\n" f"开始时间:{detail[3]}\n" f"已耗时:{detail[5]/1000}秒\n" f"扫描行数:{detail[6]}") return current_status finally: if conn: conn.close() # ------------------- 示例用法 ------------------- if __name__ == "__main__": # 替换成你的海量数据处理查询 large_data_query = """ SELECT * FROM 你的大表 WHERE 过滤条件 GROUP BY 分组列 ORDER BY 排序列 """ # 提交异步查询 query_id = submit_large_query(large_data_query) # 模拟监控逻辑:每隔10秒检查一次,直到查询完成 while True: status = track_query_status(query_id) # 查询完成的状态:SUCCESS/FAILED/CANCELLED if status in ('SUCCESS', 'FAILED', 'CANCELLED'): print(f"\n查询结束!最终状态:{status}") break print("查询仍在执行中,10秒后再次检查...\n") time.sleep(10)
注意事项
- 权限问题:确保你的角色有
MONITOR权限,或者能访问INFORMATION_SCHEMA.QUERY_HISTORY视图,不然查不到历史 - 超时设置:如果查询跑很久,要确认你的仓库没有设置过短的超时时间,必要时调整仓库的参数
- 查询ID存储:可以把查询ID存到本地文件、缓存或者数据库里,方便系统其他模块读取并监控
内容的提问来源于stack exchange,提问作者Enrai
相关产品推荐
相关产品推荐

