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

如何使用Databricks Delta Live Tables从SQL Server加载数据并传值至其他Notebook

技术方案:Delta Live Tables从SQL Server加载数据并传递至其他Notebook

一、配置SQL Server连接并通过DLT加载数据

1. 安全存储连接凭证(使用Databricks Secrets)

避免硬编码敏感信息,先在Databricks中存储SQL Server的用户名和密码:

# 通过Databricks CLI或Workspace界面创建Secret
dbutils.secrets.create(scope="sql-server-scope", key="sql-user", string_value="你的SQL Server用户名")
dbutils.secrets.create(scope="sql-server-scope", key="sql-pass", string_value="你的SQL Server密码")

2. 编写DLT管道代码(支持Python/SQL两种方式)

方式1:Python版DLT代码

创建DLT专属Notebook,实现从SQL Server拉取数据并写入Delta表:

import dlt

@dlt.table(
    name="sql_server_raw_data",
    comment="从SQL Server同步的原始业务数据",
    table_properties={
        "delta.autoOptimize.optimizeWrite": "true",
        "delta.autoOptimize.autoCompact": "true"
    }
)
def load_sql_server_data():
    # 定义SQL Server连接参数
    jdbc_url = "jdbc:sqlserver://<你的SQL Server地址>.database.windows.net:1433;databaseName=<目标数据库名>"
    conn_props = {
        "user": dbutils.secrets.get(scope="sql-server-scope", key="sql-user"),
        "password": dbutils.secrets.get(scope="sql-server-scope", key="sql-pass"),
        "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
    }
    
    # 全量读取SQL Server表(增量场景可基于时间戳/自增ID过滤)
    df = spark.read.jdbc(
        url=jdbc_url,
        table="<源表名>",
        properties=conn_props
    )
    
    return df

方式2:SQL版DLT代码

如果偏好SQL语法,可使用DLT的SQL原生支持:

CREATE OR REFRESH LIVE TABLE sql_server_raw_data
COMMENT '从SQL Server同步的原始业务数据'
TBLPROPERTIES (
    'delta.autoOptimize.optimizeWrite' = 'true',
    'delta.autoOptimize.autoCompact' = 'true'
)
AS SELECT * FROM JDBC.`jdbc:sqlserver://<你的SQL Server地址>.database.windows.net:1433;databaseName=<目标数据库名>`
OPTIONS (
    user = dbutils.secrets.get('sql-server-scope', 'sql-user'),
    password = dbutils.secrets.get('sql-server-scope', 'sql-pass'),
    driver = 'com.microsoft.sqlserver.jdbc.SQLServerDriver',
    dbtable = '<源表名>'
);

3. 部署并运行DLT管道

在Databricks Workspace中创建DLT管道,关联上述Notebook,配置合适的集群规格后启动管道,数据会自动同步到sql_server_raw_data Delta表中。

二、将数据传递至其他Notebook

根据数据量和业务场景,推荐以下3种可靠方案:

方案1:通过持久化Delta表共享数据(生产环境首选,支持大数据量)

直接在目标Notebook中读取DLT生成的Delta表即可:

# 目标Notebook代码
raw_data_df = spark.read.table("main.<你的Catalog>.<你的Schema>.sql_server_raw_data")

# 后续业务处理逻辑
display(raw_data_df)

如果需要传递特定数据子集,可在DLT中提前创建过滤视图:

# DLT Notebook中添加视图定义
@dlt.view(name="filtered_active_data")
def filter_active_records():
    return dlt.read("sql_server_raw_data").filter("status = 'active'")

目标Notebook读取视图:

active_data_df = spark.read.table("main.<你的Catalog>.<你的Schema>.filtered_active_data")

方案2:通过Notebook参数传递(适用于小数据量,如统计值、ID列表)

若仅需传递少量聚合结果或标识数据,可在DLT管道完成后触发目标Notebook并传递参数:

步骤1:在DLT中添加触发逻辑(需配置DLT管道的"Post-job tasks"或通过Job调度)

# 获取需要传递的统计数据
total_active = dlt.read("filtered_active_data").count()

# 触发目标Notebook并传递参数
dbutils.notebook.run(
    path="/Workspace/路径/到/目标Notebook",
    timeout_seconds=3600,
    arguments={
        "total_active_records": str(total_active),
        "sync_timestamp": str(spark.sql("SELECT current_timestamp()").first()[0])
    }
)

步骤2:目标Notebook接收参数

# 目标Notebook代码
total_active = dbutils.widgets.get("total_active_records")
sync_time = dbutils.widgets.get("sync_timestamp")

print(f"本次同步活跃数据量:{total_active}")
print(f"同步完成时间:{sync_time}")

方案3:通过临时视图传递(仅适用于同一会话调试,不推荐生产)

如果DLT和目标Notebook在同一交互式会话中运行,可创建临时视图共享数据:

# DLT Notebook中创建临时视图
dlt.read("sql_server_raw_data").createOrReplaceTempView("temp_sql_data")

目标Notebook读取视图:

temp_df = spark.sql("SELECT * FROM temp_sql_data")

关键注意事项

  • 确保SQL Server的1433端口对Databricks集群开放,生产环境建议使用Azure Private Link/VPC对等连接实现私有访问。
  • 增量同步场景下,建议基于SQL Server表的时间戳/自增ID字段实现增量读取,避免全量扫描:
# 增量读取示例(假设表中有update_time字段)
df = spark.read.jdbc(
    url=jdbc_url,
    table="(SELECT * FROM <源表名> WHERE update_time > '2024-05-19 00:00:00') AS incremental_data",
    properties=conn_props
)
  • 配置DLT管道权限时,需确保管道拥有Secrets读取权限和目标Catalog/Schema的写入权限。

内容的提问来源于stack exchange,提问作者bigdata techie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 12:52:41