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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 04:32:11