如何在Snowflake中读取Delta Lake最新版本并使用MERGE INTO同步数据?
问题描述
我想用MERGE INTO命令把Databricks上Delta Table的数据同步到Snowflake表中,目标是让两边记录数一致。但现在遇到问题:Delta Lake(存储在S3路径)有多个版本,Snowflake查询时会读出重复记录,请问怎么配置才能只读取Delta Lake的最新版本?
现有MERGE INTO语句
MERGE INTO myTable as target USING ( SELECT $1:DAY::TEXT AS DAY, $1:CHANNEL_CATEGORY::TEXT AS CHANNEL_CATEGORY, $1:SOURCE::TEXT AS SOURCE, $1:PLATFORM::TEXT AS PLATFROM, $1:LOB::TEXT AS LOB FROM @StageFilePathDeltaLake (FILE_FORMAT => 'sf_parquet_format') ) as src ON target.CHANNEL_CATEGORY = src.CHANNEL_CATEGORY AND target.SOURCE = src.SOURCE WHEN MATCHED THEN UPDATE SET DAY= src.DAY ,PLATFORM= src.PLATFORM ,LOB= src.LOB WHEN NOT MATCHED THEN INSERT ( DAY, CHANNEL_CATEGORY, SOURCE, PLATFORM, LOB ) VALUES ( src.DAY, src.CHANNEL_CATEGORY, src.SOURCE, src.PLATFORM, src.LOB );
现有文件格式定义
create or replace file format sf_parquet_format type = 'parquet' compression = auto;
解决方案
Delta Lake的多版本文件会导致Snowflake读取到历史数据,要只读取最新版本,可采用以下几种方式:
1. 使用Snowflake Delta专用文件格式(推荐)
Snowflake支持直接识别Delta Lake的元数据,自动过滤历史版本,只读取最新数据。只需将文件格式修改为DELTA类型:
修改文件格式
create or replace file format sf_delta_format type = 'DELTA' compression = auto;
更新MERGE INTO语句
替换文件格式为新的Delta专用格式,同时修正原语句中的拼写错误(PLATFROM改为PLATFORM):
MERGE INTO myTable as target USING ( SELECT DAY::TEXT AS DAY, CHANNEL_CATEGORY::TEXT AS CHANNEL_CATEGORY, SOURCE::TEXT AS SOURCE, PLATFORM::TEXT AS PLATFORM, LOB::TEXT AS LOB FROM @StageFilePathDeltaLake (FILE_FORMAT => 'sf_delta_format') ) as src ON target.CHANNEL_CATEGORY = src.CHANNEL_CATEGORY AND target.SOURCE = src.SOURCE WHEN MATCHED THEN UPDATE SET DAY= src.DAY, PLATFORM= src.PLATFORM, LOB= src.LOB WHEN NOT MATCHED THEN INSERT ( DAY, CHANNEL_CATEGORY, SOURCE, PLATFORM, LOB ) VALUES ( src.DAY, src.CHANNEL_CATEGORY, src.SOURCE, src.PLATFORM, src.LOB );
2. 手动解析Delta日志过滤有效文件
Delta Lake通过_delta_log目录记录版本变更,可先查询最新日志获取当前有效的数据文件列表,再在Snowflake查询中仅加载这些文件。这种方式需要手动解析日志,步骤繁琐,适合特殊场景:
- 列出S3路径下
_delta_log目录的所有文件,找到最新的.json日志文件 - 解析日志中的
add字段,提取当前有效的数据文件路径 - 在Snowflake的阶段查询中通过
PATTERN参数指定这些文件
3. 在Databricks生成最新版本快照
如果不想修改Snowflake配置,可在Databricks中将Delta Table最新版本导出为独立的Parquet文件(覆盖旧文件),再同步到S3供Snowflake读取:
# Databricks中执行,导出最新版本到指定S3路径 df = spark.read.format("delta").load("s3://your-delta-table-path") df.write.mode("overwrite").parquet("s3://your-snapshot-path")
之后将Snowflake阶段指向这个快照路径,即可读取到最新版本的数据。
内容的提问来源于stack exchange,提问作者ultraInstinct
相关产品推荐
相关产品推荐

