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

如何通过Spark Connector终止Snowflake查询或获取对应查询/会话ID

Snowflake Spark 作业关联查询终止解决方案

能否通过Spark Connector直接终止Snowflake运行中查询

目前官方Snowflake Spark Connector不支持直接主动终止正在运行的Snowflake侧查询。Spark作业被终止时,Connector默认不会主动向Snowflake发送查询取消请求,这就是你遇到的Spark作业停止后,Snowflake查询仍长时间运行的核心原因。

如何在Spark中获取Snowflake的query id、session id

你可以通过两种方式拿到所需标识,用于后续外部终止查询:

  • 方式1:获取单次查询的query id
    2.8.0及以上版本的Snowflake Spark Connector会将每次触发的查询ID存入DataFrame的查询标签中,获取代码示例:
    Scala版本:
    val df = spark.read.format("net.snowflake.spark.snowflake")
      .options(sfOptions)
      .option("query", "你的长耗时查询语句")
      .load()
    val queryId = df.queryExecution.tags.get("snowflake.query.id").orNull
    
    Python版本:
    df = spark.read.format("net.snowflake.spark.snowflake").options(**sf_options).option("query", "你的长耗时查询语句").load()
    query_id = df._jdf.queryExecution().tags().get("snowflake.query.id").get()
    
    拿到query id后可持久化存储,需要终止时在Snowflake侧执行命令 ALTER SYSTEM CANCEL QUERY '<QUERY_ID>'; 即可。
  • 方式2:固定会话ID批量管理查询
    你可以在Spark连接Snowflake的参数中手动指定session_id参数,固定本次连接使用的Snowflake会话ID,该会话下触发的所有查询都可以通过终止会话批量停止,终止会话的命令为 ALTER SYSTEM KILL SESSION '<SESSION_ID>';。

落地优化方案

你可以结合场景做如下配置,彻底解决资源浪费问题:

  • 为Spark触发的Snowflake查询设置查询超时时间:在连接参数中添加query_timeout参数,单位为秒,超过该时间的查询会被Snowflake主动终止,避免无限制跑数。
  • 给Spark作业添加终止钩子:在Spark作业中注册ShutdownHook,当作业收到终止信号时,自动调用Snowflake JDBC接口执行对应query id的取消命令,无需手动操作。
  • 若你的Spark作业运行在YARN/K8s等资源调度平台上,可以配置平台的作业终止回调逻辑,绑定对应Snowflake查询/会话的取消操作。

内容的提问来源于stack exchange,提问作者Stéphane

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 14:06:03