Azure DataBricks与Synapse Serverless动态脚本执行优化方案咨询
优化方案建议
针对你几十万条记录的脚本执行需求,结合速度和失败日志的要求,推荐以下几个更优方案:
方案1:批量分片异步执行(ADF+自定义脚本)
核心思路是把单条执行改成批量分片执行,减少ADF活动调用的开销,同时保留错误隔离和日志记录:
- 分片预处理:在SQL Server中将原表按
ViewName哈希或行号分成若干批次(比如每批次100-500条),避免单批次过大超时。 - ADF并行调度:用Lookup活动获取批次列表,For Each活动设置合理并行度(20-50,根据Synapse/ADB的并发限制调整),每个批次执行带错误处理的批量脚本。
- Synapse端脚本示例(自带日志):
-- 假设当前批次数据来自临时表或参数化查询 DECLARE @BatchScripts TABLE (ViewName NVARCHAR(255), ExternalTableScript NVARCHAR(MAX), ViewScript NVARCHAR(MAX)) INSERT INTO @BatchScripts SELECT ViewName, ExternalTableDefinition_Synapse, ViewDefinition_Synapse FROM [YourSourceTable] WHERE BatchId = @CurrentBatchId; -- 按分片逻辑筛选 DECLARE cur CURSOR FOR SELECT ViewName, ExternalTableScript, ViewScript FROM @BatchScripts DECLARE @ViewName NVARCHAR(255), @ETScript NVARCHAR(MAX), @VScript NVARCHAR(MAX) OPEN cur FETCH NEXT FROM cur INTO @ViewName, @ETScript, @VScript WHILE @@FETCH_STATUS = 0 BEGIN BEGIN TRY -- 执行外部表创建 EXEC sp_executesql @ETScript; -- 执行视图创建 EXEC sp_executesql @VScript; -- 记录成功日志 INSERT INTO Synapse_Script_Log (ViewName, Status, ExecuteTime) VALUES (@ViewName, 'Success', GETUTCDATE()); END TRY BEGIN CATCH -- 记录失败日志 INSERT INTO Synapse_Script_Log (ViewName, Status, ErrorMsg, ExecuteTime) VALUES (@ViewName, 'Failed', ERROR_MESSAGE(), GETUTCDATE()); END CATCH FETCH NEXT FROM cur INTO @ViewName, @ETScript, @VScript END CLOSE cur DEALLOCATE cur - ADB端脚本示例(PySpark批量处理):
import traceback from datetime import datetime # 读取当前批次数据(可从ADLS分片文件或直接查询SQL Server) batch_df = spark.read.jdbc( url="jdbc:sqlserver://your-sql-server:1433;databaseName=your-db;", table="(SELECT ViewName, ExternalTableDefinition_ADB, ViewDefinition_ADB FROM YourSourceTable WHERE BatchId = {}) AS BatchData".format(current_batch_id), properties={"user": "user", "password": "pwd"} ) log_records = [] for row in batch_df.collect(): view_name = row["ViewName"] et_script = row["ExternalTableDefinition_ADB"] v_script = row["ViewDefinition_ADB"] try: spark.sql(et_script) spark.sql(v_script) log_records.append((view_name, "Success", None, datetime.utcnow().isoformat())) except Exception as e: err_msg = traceback.format_exc() log_records.append((view_name, "Failed", err_msg, datetime.utcnow().isoformat())) # 写入日志到Delta表(持久化存储) log_schema = "ViewName STRING, Status STRING, ErrorMsg STRING, ExecuteTime STRING" log_df = spark.createDataFrame(log_records, schema=log_schema) log_df.write.mode("append").format("delta").save("/adls-path/adb-script-logs/") - 优势:批量减少ADF活动调用次数,并行度可控,单批次内的失败不影响其他批次,日志完整。
方案2:SQL Server生成带错误处理的脚本,直接推送到Synapse/ADB执行
跳过ADF中间环节,直接利用Synapse和ADB的批量处理能力:
- Synapse端:在SQL Server生成包含TRY/CATCH的批量脚本,用
bcp或Synapse Link推送到ADLS,然后用Synapse Serverless的OPENROWSET读取脚本文件并执行。 - ADB端:在SQL Server生成Spark SQL脚本,保存到ADLS,用ADB作业集群批量执行(作业可并行运行多个任务,每个任务处理一个脚本分片)。
- 日志处理:和方案1一致,脚本内部自带错误捕获和日志写入。
- 优势:减少ADF的调度开销,成本更低,利用原生批量执行能力提升速度。
方案3:Azure Functions+Service Bus异步并行处理
用Serverless架构实现更灵活的并行控制,避免ADF排队问题:
- 架构流程:
- 在SQL Server中分片数据,将每个批次的脚本数据序列化为JSON,发送到Azure Service Bus的两个队列(一个给Synapse,一个给ADB)。
- 创建两个Azure Functions,分别监听两个队列,触发后并行执行批次脚本。
- Functions内部捕获执行错误,将日志写入Azure SQL DB或ADLS。
- C# Function示例(Synapse处理):
using System; using System.Data.SqlClient; using System.Collections.Generic; using Newtonsoft.Json; using Microsoft.Azure.WebJobs; using Microsoft.Extensions.Logging; public static void Run([ServiceBusTrigger("synapse-script-queue", Connection = "ServiceBusConn")] string batchJson, ILogger log) { var batchItems = JsonConvert.DeserializeObject<List<ScriptItem>>(batchJson); using (var synapseConn = new SqlConnection("YourSynapseServerlessConnString")) { synapseConn.Open(); foreach (var item in batchItems) { using (var cmd = synapseConn.CreateCommand()) { try { cmd.CommandText = item.ExternalTableDefinition_Synapse; cmd.ExecuteNonQuery(); cmd.CommandText = item.ViewDefinition_Synapse; cmd.ExecuteNonQuery(); // 写入成功日志 cmd.CommandText = $"INSERT INTO ScriptLog VALUES ('{item.ViewName}', 'Success', GETUTCDATE(), NULL)"; cmd.ExecuteNonQuery(); } catch (Exception ex) { // 写入失败日志(注意转义单引号) var escapedErrMsg = ex.Message.Replace("'", "''"); cmd.CommandText = $"INSERT INTO ScriptLog VALUES ('{item.ViewName}', 'Failed', GETUTCDATE(), '{escapedErrMsg}')"; cmd.ExecuteNonQuery(); log.LogError($"Failed to process {item.ViewName}: {ex.Message}"); } } } } } public class ScriptItem { public string ViewName { get; set; } public string ExternalTableDefinition_Synapse { get; set; } public string ViewDefinition_Synapse { get; set; } } - 优势:并行控制更灵活,按执行时间计费成本更低,避免ADF的排队瓶颈。
关键优化要点
- 分片大小:根据Synapse/ADB的并发限制调整批次大小(100-500条/批次),平衡执行效率和开销。
- 限流控制:Synapse Serverless默认并发限制为20个查询,ADB集群根据资源配置调整并行任务数,避免触发限流。
- 预检查:在脚本中添加
IF NOT EXISTS判断,避免重复创建对象导致的无意义报错。 - 日志统一:将Synapse和ADB的日志统一存储到同一位置(如Azure SQL DB或ADLS Delta表),方便后续排查。
内容的提问来源于stack exchange,提问作者SouravA
相关产品推荐
相关产品推荐

