如何使用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
相关产品推荐
相关产品推荐

