如何通过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版本:
Python版本:val df = spark.read.format("net.snowflake.spark.snowflake") .options(sfOptions) .option("query", "你的长耗时查询语句") .load() val queryId = df.queryExecution.tags.get("snowflake.query.id").orNull
拿到query id后可持久化存储,需要终止时在Snowflake侧执行命令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()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
相关产品推荐
相关产品推荐

