Delta Live Tables报错处理:APPLY CHANGES源需为流式查询
问题:解决Delta Live Tables中APPLY CHANGES的流式源报错
使用场景
通过Fivetran将Oracle数据库数据同步至Databricks青铜层静态表,由Fivetran负责该表的增删改操作。
目标
借助Delta Live Tables(DLT)将数据从青铜层流式同步至银层与金层。
当前实现尝试
使用SQL笔记本编写如下代码,意图基于青铜层Delta表创建银层Delta Live Table:
CREATE OR REFRESH STREAMING TABLE cdc_test_silver; APPLY CHANGES INTO live.cdc_test_silver FROM lakehouse_poc.bronze.cdc_test KEYS (ID) SEQUENCE BY ModificationTime;
遇到的报错
Source data for the APPLY CHANGES target 'lakehouse_poc.bronze.cdc_test_silver' must be a streaming query.
解决方案
报错原因是APPLY CHANGES INTO要求源必须是流式查询,直接引用静态青铜层表属于批处理源,不符合DLT的CDC同步要求。只需将源表用STREAM()函数包装,转为流式读取即可:
修改后的代码:
CREATE OR REFRESH STREAMING TABLE cdc_test_silver; APPLY CHANGES INTO live.cdc_test_silver FROM STREAM(lakehouse_poc.bronze.cdc_test) KEYS (ID) SEQUENCE BY ModificationTime;
关键说明
STREAM()函数会将静态Delta表转为流式数据源,持续捕获Fivetran同步到青铜层的增量变更数据- 确保青铜层表为Delta格式(Fivetran同步至Databricks时默认生成Delta表)
ModificationTime字段需为递增的时间戳,用于正确排序变更事件顺序ID作为唯一键,用于匹配银层表中的行以执行增删改操作
内容的提问来源于stack exchange,提问作者play_something_good
相关产品推荐
相关产品推荐

