如何通过API或类库读取Azure Databricks中SQL命令的输出供.NET Core使用
实现方案:从.NET Core调用Azure Databricks Notebook并获取SQL执行结果
前置准备
- 首先获取Azure Databricks工作区的访问令牌,可从Databricks用户设置界面生成,令牌需妥善保管用于后续接口鉴权
- 记录目标Notebook的存储路径、Databricks工作区的区域域名(格式为
https://<你的区域>.azuredatabricks.net) - .NET Core项目建议使用.NET 6及以上版本,API兼容性更好
方案一:调用Databricks Jobs API触发Notebook运行并拉取结果
这是适配性最高的通用方案,操作步骤如下:
- 先在Databricks Notebook中添加结果返回逻辑,SQL执行完成后把结果转成结构化格式,用
dbutils.notebook.exit()返回,示例Notebook代码:
-- 你的业务SQL执行逻辑 SELECT * FROM your_table WHERE condition = 'xxx';
# 捕获上面SQL单元格的执行结果转为JSON格式返回 sql_result = _sqldf.toJSON().collect() dbutils.notebook.exit(sql_result)
注意:
_sqldf是Databricks SQL单元格执行后默认生成的DataFrame变量名,如果你是用Python调用spark.sql执行SQL,可以自行赋值变量名替换。
- .NET Core侧直接调用Databricks REST API实现任务提交、状态轮询、结果拉取,示例代码:
using System.Net.Http.Headers; using System.Text.Json; // 基础配置参数,建议从配置中心/环境变量读取,不要硬编码 const string DATABRICKS_DOMAIN = "https://<你的区域>.azuredatabricks.net"; const string ACCESS_TOKEN = "<你的Databricks访问令牌>"; const string NOTEBOOK_PATH = "/Users/xxx/your_notebook_name"; var httpClient = new HttpClient(); httpClient.DefaultRequestHeaders.Authorization = new AuthenticationHeaderValue("Bearer", ACCESS_TOKEN); // 1. 提交Notebook运行任务 var runPayload = new { run_name = "dotnet-call-sql-run", notebook_task = new { notebook_path = NOTEBOOK_PATH, // 可选:传递动态参数给Notebook,用于SQL查询条件传参 base_parameters = new { filter = "xxx" } }, // 可替换为已有集群ID,避免每次新建集群浪费资源 new_cluster = new { spark_version = "13.3.x-scala2.12", node_type_id = "Standard_DS3_v2", num_workers = 1 } }; var runResponse = await httpClient.PostAsJsonAsync($"{DATABRICKS_DOMAIN}/api/2.1/jobs/runs/submit", runPayload); runResponse.EnsureSuccessStatusCode(); var runResult = await runResponse.Content.ReadFromJsonAsync<JsonElement>(); var runId = runResult.GetProperty("run_id").GetInt64(); // 2. 轮询任务运行状态直到结束 string runState = ""; JsonElement sqlOutput = default; while (runState != "TERMINATED" && runState != "FAILED" && runState != "SKIPPED") { await Task.Delay(5000); // 轮询间隔可自行调整,避免触发API限流 var statusResponse = await httpClient.GetAsync($"{DATABRICKS_DOMAIN}/api/2.1/jobs/runs/get?run_id={runId}"); statusResponse.EnsureSuccessStatusCode(); var statusResult = await statusResponse.Content.ReadFromJsonAsync<JsonElement>(); runState = statusResult.GetProperty("state").GetProperty("life_cycle_state").GetString(); if (runState == "TERMINATED") { // 3. 拉取Notebook返回的SQL结果 var outputResponse = await httpClient.GetAsync($"{DATABRICKS_DOMAIN}/api/2.1/jobs/runs/get-output?run_id={runId}"); outputResponse.EnsureSuccessStatusCode(); var outputResult = await outputResponse.Content.ReadFromJsonAsync<JsonElement>(); sqlOutput = outputResult.GetProperty("notebook_output").GetProperty("result"); break; } } // 最终sqlOutput就是结构化的SQL执行结果,可反序列化为自定义模型使用 Console.WriteLine(sqlOutput.ToString());
方案二:直接调用Databricks SQL Warehouse执行SQL(无需提前维护Notebook)
如果你的需求仅为执行指定SQL,不需要用到Notebook内的其他逻辑,可以直接使用SQL Warehouse API,流程更简单:
- 先在Databricks工作区创建SQL Warehouse,记录它的Warehouse ID
- .NET Core侧直接调用SQL执行接口即可,示例代码:
// HttpClient鉴权逻辑和方案一完全一致 var sqlPayload = new { warehouse_id = "<你的SQL Warehouse ID>", statement = "SELECT * FROM your_table WHERE condition = 'xxx'", wait_timeout = "30s" }; var sqlResponse = await httpClient.PostAsJsonAsync($"{DATABRICKS_DOMAIN}/api/2.0/sql/statements", sqlPayload); sqlResponse.EnsureSuccessStatusCode(); var sqlResult = await sqlResponse.Content.ReadFromJsonAsync<JsonElement>(); // 直接解析返回的结果集即可 var resultRows = sqlResult.GetProperty("result").GetProperty("data_array");
注意事项
- Databricks访问令牌禁止硬编码在代码里,建议使用.NET Core配置中心、Azure Key Vault等加密存储
- 轮询间隔不要太短,避免触发API限流规则
- 如果SQL返回结果集超过10MB,建议不要直接通过接口返回,先将结果写入Azure Blob Storage/ADLS,.NET Core侧再从存储拉取文件,避免接口响应超限
内容的提问来源于stack exchange,提问作者Prashant Rewatkar
相关产品推荐
相关产品推荐

