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

Azure Data Factory数据流:查询模式下实现CosmosDB变更馈送的替代方案

解决方案

一、关于你当前Alter Row的使用逻辑

你用Alter row修饰符将Upsert设为true()是正确的第一步,但这只是给数据流中的行打上Upsert标记,要让Delta表真正执行Upsert,还需在Sink(Delta表)的设置中完成对应配置:

  • 在Sink的Update method中选择Upsert
  • 指定Key columns(用于匹配Delta表已有数据的主键列,比如CosmosDB的id字段)
  • 若需处理删除操作,要在Alter row里新增删除标记逻辑:比如当数据中的删除标识字段(如isDeleted)为true()时,将Delete设为true(),同时在Sink的Update method中勾选Allow delete

二、替代CosmosDB变更馈送的增量方案(适配Query作为源的场景)

由于必须用Query作为输入,无法依赖ADF内置的变更馈送,需自行实现增量拉取逻辑,核心是跟踪上次同步的时间戳,仅拉取新增/变更的数据:

1. 存储同步Checkpoint

在ADF中用变量或**独立存储表(如另一个Delta表、Azure Blob文件)**记录每次同步的最后时间戳(可使用CosmosDB文档自带的_ts字段,这是文档最后修改的Unix时间戳)。

2. 编写增量Query

每次同步时,从Checkpoint读取上次的时间戳,编写参数化的CosmosDB Query拉取_ts大于该时间戳的文档:

SELECT * FROM c WHERE c._ts > @lastSyncTimestamp

在ADF数据流的源配置中,将该Query作为输入,通过参数传递@lastSyncTimestamp的值。

3. 处理删除操作

  • 若为软删除(文档保留但标记删除状态):直接在Query中包含这些文档,通过Alter row根据删除标记设置Delete = true()
  • 若为硬删除:需额外维护一张删除日志表记录被删文档的ID,同步时拉取这些ID,在数据流中生成对应删除行并标记Delete = true()

4. 更新Checkpoint

每次同步完成后,将本次同步的最大_ts值更新到Checkpoint中,确保下一次仅拉取新增的变更数据。

三、完整流程总结

  1. 读取上次同步的时间戳Checkpoint
  2. 通过参数化Query从CosmosDB拉取_ts大于该时间戳的所有变更文档(含新增、更新、软删除文档)
  3. 使用Alter row修饰符:
    • 对需更新/插入的行,设置Upsert = true()
    • 对需删除的行,设置Delete = true()
  4. 将数据写入Delta表Sink,开启Upsert模式并指定主键,同时勾选Allow delete
  5. 更新Checkpoint为本次同步的最大_ts值

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 14:43:16