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锁查看权限,现提出以下问题:
- 如何通过PySpark代码或UI终止此类进程?
- 如何优化上述代码以提升性能与容错性,尤其是使用
bulkCopyTableLock、bulkCopyTimeout、分区等参数? - 有哪些资源可帮助我更好地理解此类问题?
执行的代码如下:
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
相关产品推荐
相关产品推荐

