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

如何通过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运行并拉取结果

这是适配性最高的通用方案,操作步骤如下:

  1. 先在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,可以自行赋值变量名替换。

  1. .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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 15:09:03