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

Databricks DLT与CDC:底层数据变更时APPLY CHANGES报错及解决方案咨询

处理S3 Parquet CDC场景的DLT构建方案

核心问题说明

流表搭配APPLY CHANGES报错的原因是:S3 Parquet的文件修改操作不属于追加型流源的追踪范围,流模式仅能识别新增文件,无法捕获现有文件的内容变更。因此需要切换到「全量+增量结合的Live表」模式来覆盖新增和更新场景。


Bronze层实现:捕获所有新增/变更数据

放弃流表,改用普通Live表,通过文件元数据过滤实现增量扫描,避免每次全量读取:

CREATE LIVE TABLE bronze_cdc_data
COMMENT "Bronze层存储S3 Parquet的全量CDC数据(含新增、变更)"
AS SELECT 
  *,
  _metadata.file_path AS source_file_path,
  _metadata.file_modification_time AS source_modified_time
FROM cloud_files(
  "s3://your-bucket/target-path/",
  "parquet"
)
-- 仅加载上次运行后新增/修改的文件
WHERE source_modified_time > (
  SELECT COALESCE(MAX(source_modified_time), '1970-01-01') 
  FROM LIVE.bronze_cdc_data
);
  • 逻辑说明:利用S3文件的file_modification_time元数据,每次运行只读取上次同步后更新的文件,既覆盖新增文件,也能捕获现有文件的修改(文件修改后file_modification_time会更新)
  • 适配场景:支持S3上的文件追加、覆盖修改、分区文件更新等操作

Silver层实现:合并CDC数据到目标表

基于Bronze层的增量数据,用APPLY CHANGES完成新增/更新的合并,分两种场景处理:

场景1:源数据带CDC操作标识(如op_type)

如果源数据包含明确的操作类型(INSERT/UPDATE/DELETE)和主键:

-- 初始化Silver层目标表
CREATE LIVE TABLE silver_cdc_target
COMMENT "Silver层最终合并表"
TBLPROPERTIES ("delta.enableChangeDataFeed" = "true")
AS SELECT id, col1, col2, op_type, update_ts
FROM LIVE.bronze_cdc_data
WHERE 1=0;

-- 应用CDC变更
APPLY CHANGES INTO LIVE.silver_cdc_target
FROM LIVE.bronze_cdc_data
KEYS (id) -- 主键字段
APPLY AS DELETE WHEN op_type = 'DELETE'
APPLY AS UPDATE WHEN op_type = 'UPDATE'
APPLY AS INSERT WHEN op_type = 'INSERT'
SEQUENCE BY update_ts; -- 排序字段,确保变更顺序

场景2:源数据无CDC标识,仅按主键保留最新记录

如果源数据只有主键和更新时间,先在Bronze层去重,再合并到Silver层:

-- 中间表:Bronze层去重,保留每个主键的最新记录
CREATE LIVE TABLE bronze_deduped
COMMENT "Bronze层去重后的最新CDC记录"
AS SELECT id, col1, col2, update_ts
FROM (
  SELECT 
    *,
    ROW_NUMBER() OVER (PARTITION BY id ORDER BY update_ts DESC) AS rn
  FROM LIVE.bronze_cdc_data
)
WHERE rn = 1;

-- 合并到Silver层目标表
APPLY CHANGES INTO LIVE.silver_cdc_target
FROM LIVE.bronze_deduped
KEYS (id)
SEQUENCE BY update_ts
WHEN NOT MATCHED THEN INSERT *
WHEN MATCHED THEN UPDATE SET *;

关键注意事项

  • 禁止用流表处理S3上有文件修改的场景,流模式仅支持追踪新增文件
  • Bronze层的过滤逻辑依赖S3文件的file_modification_time,如果是手动修改文件需确保该元数据正常更新
  • Silver层使用APPLY CHANGES时,必须指定主键和排序字段,避免数据冲突
  • 若S3存在大量小文件,建议定期合并小文件,提升DLT pipeline运行效率

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 01:55:03