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

通过ADF实现CDC后,将Blob日志同步至Databricks Delta表求助

Azure SQL CDC日志同步至Databricks Delta表解决方案

1. 优化Blob中CDC文件的处理(可选)

ADF CDC默认生成独立CSV文件,可通过以下方式减少文件数量:

  • 在ADF的CDC Sink配置中,启用合并文件选项,设置按时间窗口(如小时)合并输出文件,写入指定目录(如/processed_cdc/)
  • 若不调整ADF,也可直接在Databricks中读取整个CDC目录下的所有CSV文件,后续处理逻辑兼容多文件场景

2. 实时同步至Delta表(推荐用Auto Loader)

利用Databricks Auto Loader监听Blob目录的新增CDC文件,自动执行Upsert/Delete操作更新Delta表,代码示例:

# 替换为你的Blob存储路径
cdc_source_path = "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/cdc_logs/"
# 替换为你的Delta表存储路径
delta_target_path = "/delta/cdc_target_table"

# 用Auto Loader增量读取CDC CSV,自动处理Schema变更
raw_cdc_df = spark.readStream \
    .format("cloudFiles") \
    .option("cloudFiles.format", "csv") \
    .option("cloudFiles.schemaLocation", "/dbfs/tmp/cdc_schema_store") \
    .option("header", "true") \
    .load(cdc_source_path)

# 过滤Azure SQL CDC的无效记录(保留插入/更新后/删除操作)
filtered_cdc_df = raw_cdc_df.filter("__$operation IN (2,4,1)") \
    .withColumn("operation_type", 
        when(col("__$operation") == 2, "I")
        .when(col("__$operation") == 4, "U")
        .when(col("__$operation") == 1, "D")
    )

# 定义Upsert逻辑
from delta.tables import DeltaTable

def apply_cdc_changes(micro_batch_df, batch_id):
    delta_table = DeltaTable.forPath(spark, delta_target_path)
    
    delta_table.alias("target") \
        .merge(
            micro_batch_df.alias("source"),
            "target.your_primary_key = source.your_primary_key"  # 替换为实际主键字段
        ) \
        .whenMatchedDelete(condition="source.operation_type = 'D'") \
        .whenMatchedUpdate(set={
            "column1": "source.column1",
            "column2": "source.column2",
            "last_updated": "source.__$start_lsn"  # 用CDC的LSN作为更新时间
        }) \
        .whenNotMatchedInsert(condition="source.operation_type != 'D'", values={
            "your_primary_key": "source.your_primary_key",
            "column1": "source.column1",
            "column2": "source.column2",
            "created_at": "source.__$start_lsn"
        }) \
        .execute()

# 启动流处理,启用检查点实现断点续传
filtered_cdc_df.writeStream \
    .format("delta") \
    .foreachBatch(apply_cdc_changes) \
    .option("checkpointLocation", "/dbfs/tmp/cdc_checkpoint") \
    .start()

3. 批量同步方案(分钟级延迟)

若不需要严格实时,可通过ADF定时触发Databricks任务批量处理:

  1. 在ADF创建时间触发管道,设置触发间隔(如5分钟)
  2. 管道中添加Databricks Notebook Activity,执行以下批量处理代码:
# 读取未处理的CDC文件
cdc_df = spark.read \
    .option("header", "true") \
    .csv("abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/cdc_logs/")

# 过滤并转换CDC操作类型
filtered_cdc_df = cdc_df.filter("__$operation IN (2,4,1)") \
    .withColumn("operation_type", 
        when(col("__$operation") == 2, "I")
        .when(col("__$operation") == 4, "U")
        .when(col("__$operation") == 1, "D")
    )

# 执行Upsert到Delta表
delta_table = DeltaTable.forPath(spark, "/delta/cdc_target_table")
delta_table.alias("target").merge(
    filtered_cdc_df.alias("source"),
    "target.your_primary_key = source.your_primary_key"
) \
.whenMatchedDelete(condition="source.operation_type = 'D'") \
.whenMatchedUpdate(set={"column1": "source.column1", "column2": "source.column2"}) \
.whenNotMatchedInsert(condition="source.operation_type != 'D'", values={"your_primary_key": "source.your_primary_key", "column1": "source.column1", "column2": "source.column2"}) \
.execute()

# 归档已处理文件,避免重复执行
dbutils.fs.mv(
    "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/cdc_logs/*.csv",
    "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/cdc_archived/",
    recurse=True
)

关键注意事项

  • 主键约束:必须确保CDC日志包含源表的主键字段,否则无法准确匹配数据行
  • CDC操作过滤:Azure SQL CDC的__$operation字段中,3代表更新前的旧数据,需过滤掉避免重复处理
  • Schema兼容性:若源表Schema变更,Auto Loader会自动同步Schema,Delta表需设置spark.databricks.delta.schema.autoMerge.enabled = true开启自动Schema合并
  • 检查点与归档:流处理的检查点需存储在持久化路径,批量处理必须归档已处理文件,防止重复同步

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 04:52:42