通过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任务批量处理:
- 在ADF创建时间触发管道,设置触发间隔(如5分钟)
- 管道中添加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
相关产品推荐
相关产品推荐

