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

如何将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的增删改操作:

  1. 开启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;
    
  2. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 02:30:40