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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:18:43