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

Spark UI无法终止运行查询且写入Azure SQL Server失败的技术问询

问题描述

在Azure Databricks中执行以下PySpark代码,读取Hive表Mechanics全量数据写入SQL Server时,进程持续运行数分钟无法终止,Spark UI仅显示「Running Queries(1)」且无终止选项,也未生成任何Job或Stage。

替换目标表为其他表或glb.mechanics_invoices_temp时,代码可正常运行。

Hive表Mechanics包含2700490条数据,无SQL Server锁查看权限,现提出以下问题:

  1. 如何通过PySpark代码或UI终止此类进程?
  2. 如何优化上述代码以提升性能与容错性,尤其是使用bulkCopyTableLock、bulkCopyTimeout、分区等参数?
  3. 有哪些资源可帮助我更好地理解此类问题?

执行的代码如下:

df = spark.sql("select * from Mechanics")
df.write.format("com.microsoft.sqlserver.jdbc.spark")\
.mode("overwrite").option("url",sql_connection_properties)\
.option("dbtable",glb.mechanics_invoices)\
.option("user",username).option("password",pwd)\
.option("bulkCopyBatchSize",10000).save()

问题解答

1. 终止进程的方法

通过Databricks UI操作

  • 打开对应笔记本,点击右上角「Run」下拉菜单,选择「Cancel Runs」终止当前所有运行命令。
  • 进入Compute页面,找到运行代码的集群,点击「Terminate」直接终止集群(注意:会中断集群上所有任务,需谨慎操作)。
  • 进入Spark UI的「SQL」标签页,找到对应Running Query,点击「Kill Query」按钮;若未显示该按钮,可查看集群的「Driver Logs」获取进程ID,联系管理员协助终止。

通过PySpark代码终止

  • 在同一笔记本中执行以下代码,终止当前会话的所有任务:
# 取消当前会话的所有Job
spark.sparkContext.cancelAllJobs()
  • 若上述方法无效,可使用Databricks REST API终止指定运行任务(需配置认证信息):
import requests
import json

headers = {
    "Authorization": "Bearer <你的认证Token>",
    "Content-Type": "application/json"
}
# 替换为你的工作区URL和运行任务ID
response = requests.post(
    "https://<你的工作区URL>/api/2.0/jobs/runs/cancel",
    headers=headers,
    data=json.dumps({"run_id": "<目标运行任务ID>"})
)

2. 代码优化建议

核心参数调整

  • bulkCopyTableLock:默认值为true,批量复制期间会锁定目标表,避免外部操作干扰;若目标表有其他业务访问,可设置为false,但可能增加冲突风险。示例:
    .option("bulkCopyTableLock", "true")
    
  • bulkCopyTimeout:设置超时时间(单位:秒),避免进程无限期挂起,建议根据数据量设置为300-1800秒。示例:
    .option("bulkCopyTimeout", "600")
    
  • bulkCopyBatchSize:当前设置为10000,可根据SQL Server性能调整,比如增大到20000或减小到5000,找到最优批量大小。

分区与并行优化

  • 利用Hive表分区:如果Mechanics表本身有分区,读取时指定分区过滤减少数据量;若无分区,读取后对DataFrame重分区,提升写入并行度(建议设置为集群核心数的2-4倍):
    df = spark.sql("select * from Mechanics").repartition(16)
    
  • 调整Spark并行任务数:设置spark.sql.shuffle.partitions参数,匹配集群资源:
    spark.conf.set("spark.sql.shuffle.partitions", "16")
    

容错性优化

  • 原子替换目标表:先写入临时表,再通过SQL Server的原子操作替换目标表,减少锁持有时间:
    # 先写入临时表
    df.write.format("com.microsoft.sqlserver.jdbc.spark")\
    .mode("overwrite").option("url",sql_connection_properties)\
    .option("dbtable",glb.mechanics_invoices_temp)\
    .option("user",username).option("password",pwd)\
    .option("bulkCopyBatchSize",10000).save()
    
    # 执行SQL Server原子替换
    conn = spark._sc._gateway.jvm.java.sql.DriverManager.getConnection(
        sql_connection_properties, username, pwd
    )
    stmt = conn.createStatement()
    stmt.execute("DROP TABLE IF EXISTS glb.mechanics_invoices; EXEC sp_rename 'glb.mechanics_invoices_temp', 'mechanics_invoices';")
    stmt.close()
    conn.close()
    
  • 增加重试机制:使用重试库处理临时连接或锁问题:
    from tenacity import retry, stop_after_attempt, wait_exponential
    
    @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10))
    def write_to_sqlserver(df):
        df.write.format("com.microsoft.sqlserver.jdbc.spark")\
        .mode("overwrite").option("url",sql_connection_properties)\
        .option("dbtable",glb.mechanics_invoices)\
        .option("user",username).option("password",pwd)\
        .option("bulkCopyBatchSize",10000)\
        .option("bulkCopyTimeout", "600")\
        .save()
    
    write_to_sqlserver(df)
    

3. 学习资源推荐

  • Databricks官方文档:重点查看「Spark SQL Server Connector」章节,了解参数配置、性能调优和故障排查方法。
  • Spark官方文档:学习Spark JDBC写入的原理、并行度调整和故障处理机制。
  • SQL Server官方文档:了解批量插入(Bulk Copy)的锁机制、性能优化建议,以及表锁排查方法。
  • Databricks社区论坛:搜索类似的SQL Server写入挂起问题,参考其他用户的解决方案和经验。

内容的提问来源于stack exchange,提问作者Ramaraju.d

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 22:35:18