Azure Databricks增量加载中跟踪Snowflake源删除记录的有效方法
解决方案:Snowflake只读数据源同步至Databricks(含增量更新+删除同步)
针对你提到的Snowflake只读数据源、已实现时间戳增量加载、需同步删除且不能定期全量的场景,以下是实用的落地方案:
核心策略选择
1. 优先利用软删除标记(推荐)
如果Snowflake源表本身带有软删除标识(比如is_deleted布尔字段、deleted_at时间戳字段),这是最高效的方式——直接通过增量拉取包含删除标记的记录,在Databricks端用MERGE语句完成更新+删除同步。
2. 主键差异对比(无软删除标记时)
若源表没有软删除机制,只能通过对比源端与目标端的主键来识别删除记录,但要通过时间分片缩小对比范围避免全量扫描的低效问题。
具体实现示例
场景1:源表含软删除标记(如deleted_at)
假设你已在Databricks中维护了一张同步元数据表sync_metadata,用于存储每张表的上次同步最大时间戳。
步骤1:获取上次同步时间戳
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 读取元数据,获取目标表的上次同步时间 last_sync_ts = spark.sql(""" SELECT max_sync_ts FROM sync_metadata WHERE table_name = 'snowflake_customer_table' """).collect()[0][0]
步骤2:增量拉取Snowflake数据
拉取所有在上次同步后更新或删除的记录:
snowflake_inc_df = spark.read \ .format("snowflake") \ .option("sfUrl", "your-account.snowflakecomputing.com") \ .option("sfUser", "your-user") \ .option("sfPassword", "your-password") \ .option("sfDatabase", "your-db") \ .option("sfSchema", "your-schema") \ .option("sfWarehouse", "your-warehouse") \ .option("query", f""" SELECT customer_id, name, email, updated_at, deleted_at FROM customer_table WHERE updated_at > '{last_sync_ts}' OR deleted_at > '{last_sync_ts}' """) \ .load()
步骤3:合并更新至Databricks目标表
用MERGE语句完成插入、更新、删除的原子操作:
snowflake_inc_df.createOrReplaceTempView("inc_updates") spark.sql(""" MERGE INTO your_target_db.customer_target t USING inc_updates u ON t.customer_id = u.customer_id -- 匹配到且已删除的记录,执行删除 WHEN MATCHED AND u.deleted_at IS NOT NULL THEN DELETE -- 匹配到未删除的记录,执行更新 WHEN MATCHED THEN UPDATE SET t.name = u.name, t.email = u.email, t.updated_at = u.updated_at -- 未匹配到且未删除的记录,执行插入 WHEN NOT MATCHED AND u.deleted_at IS NULL THEN INSERT (customer_id, name, email, updated_at) VALUES (u.customer_id, u.name, u.email, u.updated_at) """) # 更新同步元数据,记录本次同步的最大时间戳 new_max_ts = spark.sql(""" SELECT GREATEST(COALESCE(max(updated_at), '1970-01-01'), COALESCE(max(deleted_at), '1970-01-01')) FROM inc_updates """).collect()[0][0] spark.sql(f""" UPDATE sync_metadata SET max_sync_ts = '{new_max_ts}' WHERE table_name = 'snowflake_customer_table' """)
场景2:源表无软删除标记,用主键差异同步删除
这种方式需要将增量更新和删除同步拆分执行,重点是缩小主键对比的范围,避免全量扫描。
步骤1:常规增量更新(同之前逻辑)
基于updated_at拉取增量数据,用MERGE完成插入和更新,更新max_sync_ts。
步骤2:定期同步删除(例如每天执行一次)
通过对比最近时间窗口内的主键,识别已删除的记录:
# 读取上次删除同步的时间戳(需在sync_metadata中新增该字段) last_delete_sync_ts = spark.sql(""" SELECT last_delete_sync_ts FROM sync_metadata WHERE table_name = 'snowflake_customer_table' """).collect()[0][0] # 拉取Snowflake源表中最近更新/创建的主键 source_keys = spark.read \ .format("snowflake") \ .option("<snowflake-configs>") \ .option("query", f""" SELECT customer_id FROM customer_table WHERE updated_at >= '{last_delete_sync_ts}' UNION SELECT customer_id FROM customer_table WHERE created_at >= '{last_delete_sync_ts}' """) \ .load() # 拉取Databricks目标表中对应时间范围的主键 target_keys = spark.sql(f""" SELECT customer_id FROM your_target_db.customer_target WHERE updated_at >= '{last_delete_sync_ts}' OR created_at >= '{last_delete_sync_ts}' """) # 找出目标表存在但源表不存在的主键(即已删除的记录) deleted_keys = target_keys.join(source_keys, on="customer_id", how="left_anti") deleted_keys.createOrReplaceTempView("deleted_customer_ids") # 删除目标表中的对应记录 spark.sql(""" DELETE FROM your_target_db.customer_target WHERE customer_id IN (SELECT customer_id FROM deleted_customer_ids) """) # 更新删除同步的时间戳 spark.sql(f""" UPDATE sync_metadata SET last_delete_sync_ts = CURRENT_TIMESTAMP() WHERE table_name = 'snowflake_customer_table' """)
关键优化建议
- 时间分片策略:无软删除时,不要全量对比所有主键,每次只对比最近3-7天的数据,每周再执行一次范围更大的对比(比如最近30天),平衡效率和数据一致性。
- Snowflake查询优化:确保源表的
updated_at、created_at字段有索引,减少主键拉取的查询时间;利用Snowflake的分区表特性,只扫描有变更的分区。 - 元数据原子性:更新
sync_metadata时,尽量用事务包裹(Databricks Delta Lake支持ACID事务),避免同步失败导致的元数据与实际数据不一致。
内容的提问来源于stack exchange,提问作者Saswat Ray
相关产品推荐
相关产品推荐

