如何将Azure SQL DB数据加载至Synapse Spark DataLake并同步更新?
Azure SQL DB 与 Synapse Spark DataLake 数据同步解决方案
一、CSV数据导入Azure SQL DB
先把ADLS中的CSV数据导入Azure SQL DB,后续直接在SQL库内完成增改操作:
- 方式1:PySpark批量写入
适合结构化CSV数据,通过PySpark读取ADLS文件后写入SQL DB:# 读取ADLS中的CSV文件 df = spark.read.csv("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/<CSV路径>", header=True, inferSchema=True) # 写入Azure SQL DB df.write \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<SQL服务器名>.database.windows.net:1433;databaseName=<数据库名>;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;") \ .option("dbtable", "<目标SQL表名>") \ .option("user", "<SQL用户名>") \ .option("password", "<SQL密码>") \ .mode("append") \ # 全量覆盖用"overwrite",追加用"append" .save() - 方式2:SQL Bulk Insert(大文件首选)
先将CSV所在的ADLS路径挂载为SQL DB外部数据源,再执行批量插入:-- 创建外部数据源 CREATE EXTERNAL DATA SOURCE ADLS_CSV_Source WITH ( TYPE = BLOB_STORAGE, LOCATION = 'https://<存储账户名>.blob.core.windows.net/<容器名>', CREDENTIAL = <存储账户SAS凭证> ); -- 批量插入数据 BULK INSERT <目标SQL表名> FROM '<CSV文件路径>' WITH ( DATA_SOURCE = 'ADLS_CSV_Source', FORMAT = 'CSV', FIRSTROW = 2, -- 跳过表头行 FIELDTERMINATOR = ',', ROWTERMINATOR = '\n' );
二、PySpark连接Azure SQL DB与Spark DataLake
在Synapse Notebook中实现双向数据交互:
1. 从Azure SQL DB读取数据到Spark DataFrame
# 读取SQL DB指定表的数据 sql_df = spark.read \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<SQL服务器名>.database.windows.net:1433;databaseName=<数据库名>;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;") \ .option("dbtable", "<SQL表名>") \ .option("user", "<SQL用户名>") \ .option("password", "<SQL密码>") \ .load()
2. 将处理后的数据写入Spark DataLake(注册湖表)
把SQL DB读取的数据转换后,写入ADLS并注册为Spark湖表:
# 示例:数据转换(新增计算列) transformed_df = sql_df.withColumn("计算列", sql_df["原字段"] * 2) # 写入ADLS(用Delta Lake支持ACID和增量更新) transformed_df.write \ .format("delta") \ .mode("overwrite") \ # 增量追加用"append" .save("abfss://<容器名>@<存储账户名>.dfs.core.windows.net<湖表存储路径>") # 注册为湖数据库中的Spark表 spark.sql(""" CREATE OR REPLACE TABLE <湖数据库名>.<Spark表名> USING DELTA LOCATION 'abfss://<容器名>@<存储账户名>.dfs.core.windows.net<湖表存储路径>' """)
三、SQL DB数据增改后同步至Spark DataLake
确保SQL DB的变更能同步到Spark湖表,推荐两种方案:
方案1:基于Azure SQL DB CDC的增量同步
适合准实时同步场景,捕获SQL DB的增删改操作:
- 开启SQL DB的CDC功能:
-- 启用数据库级CDC EXEC sys.sp_cdc_enable_db; -- 启用目标表的CDC EXEC sys.sp_cdc_enable_table @source_schema = N'<表所属Schema>', @source_name = N'<目标表名>', @role_name = NULL, @supports_net_changes = 1; - 在Synapse Pipeline中定时触发Notebook,合并变更数据到湖表:
# 读取SQL DB的CDC变更表 cdc_df = spark.read \ .format("jdbc") \ .option("url", "<SQL DB连接地址>") \ .option("dbtable", "cdc.<Schema>_<表名>_CT") \ .option("user", "<SQL用户名>") \ .option("password", "<SQL密码>") \ .load() # 合并到Delta湖表 cdc_df.createOrReplaceTempView("cdc_changes") spark.sql(""" MERGE INTO <湖数据库名>.<Spark表名> target USING cdc_changes source ON target.id = source.id WHEN MATCHED AND source.__$operation IN (1, 4) THEN DELETE WHEN MATCHED AND source.__$operation IN (2, 3) THEN UPDATE SET * WHEN NOT MATCHED AND source.__$operation IN (2, 4) THEN INSERT * """)
方案2:定时全量/增量刷新(小数据量场景)
通过Synapse Pipeline定时触发Notebook,执行同步:
- 全量刷新:直接读取SQL DB全表覆盖湖表
- 增量刷新:基于时间戳/自增ID读取最新数据,合并到湖表
# 读取SQL DB中最近24小时的增量数据 incremental_df = spark.read \ .format("jdbc") \ .option("url", "<SQL DB连接地址>") \ .option("dbtable", "(SELECT * FROM <SQL表名> WHERE 更新时间 >= DATEADD(HOUR, -24, GETDATE())) AS incremental_data") \ .option("user", "<SQL用户名>") \ .option("password", "<SQL密码>") \ .load() # 合并到Delta湖表 incremental_df.createOrReplaceTempView("incremental_data") spark.sql(""" MERGE INTO <湖数据库名>.<Spark表名> target USING incremental_data source ON target.id = source.id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * """)
四、验证Notebook读取最新数据
在Synapse Notebook中读取湖表,确认获取最新同步数据:
# 读取湖数据库中的Spark表 latest_df = spark.read.table("<湖数据库名>.<Spark表名>") # 查看数据验证 latest_df.show(10)
用Delta Lake的话,还可以查看表版本历史确认同步效果:
DESCRIBE HISTORY <湖数据库名>.<Spark表名>
内容的提问来源于stack exchange,提问作者BigData Lover
相关产品推荐
相关产品推荐

