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

Azure Synapse翻滚窗口触发器增量加载机制及自定义逻辑问询

Azure Synapse翻滚窗口触发器与增量加载说明

1. 触发器的核心能力边界

Azure Synapse的翻滚窗口触发器不具备自动识别增量数据的能力,它的核心作用只是按预设的固定时间窗口(比如每小时、每天)定时触发流水线运行,不会主动跟踪数据源的生成/修改时间,也无法自动筛选出自上次运行以来的增量数据。触发后流水线是全量加载还是增量加载,完全由你编写的业务逻辑决定,触发器只负责“按时启动”这一步。

2. 默认加载行为:无内置增量逻辑

如果不额外配置自定义逻辑,流水线触发后会执行你编写的全量数据处理逻辑——比如直接读取整个数据源的所有数据,不会自动过滤增量内容。

3. 自定义增量加载的实操方法

方法1:利用触发器内置窗口参数

翻滚窗口触发器会自动生成windowStartTime和windowEndTime两个系统参数,你可以直接在流水线中引用这两个参数,匹配数据的时间范围实现增量筛选:

  • SQL场景示例:在数据读取的SQL语句中嵌入参数
    SELECT * FROM source_table 
    WHERE create_time >= '@{pipeline().parameters.windowStartTime}' 
      AND create_time < '@{pipeline().parameters.windowEndTime}'
    
  • 注意:需确保数据源有可靠的时间戳列(如创建时间、修改时间)来对齐窗口范围。

方法2:用水印表追踪增量

如果数据时间范围和窗口不完全对齐,或需要跟踪修改过的数据,可以用专门的水印表存储上次处理的临界点:

  1. 创建水印表(示例结构):
    CREATE TABLE watermark_table (
        pipeline_name VARCHAR(100) PRIMARY KEY,
        last_processed_time DATETIME
    )
    
  2. 流水线启动时,先读取水印表中的last_processed_time作为增量起始点
  3. 筛选增量数据:
    SELECT * FROM source_table 
    WHERE modify_time > (SELECT last_processed_time FROM watermark_table WHERE pipeline_name = 'your_pipeline')
      AND modify_time <= '@{pipeline().parameters.windowEndTime}'
    
  4. 数据处理完成后,更新水印表的last_processed_time为当前窗口的windowEndTime

方法3:基于存储系统的文件增量筛选

如果数据源是Azure Blob Storage、ADLS Gen2这类存储,可以通过文件的最后修改时间筛选增量文件:

  • 用Get Metadata活动获取存储路径下的所有文件列表
  • 用Filter活动过滤出最后修改时间大于上次运行时间的文件
  • 将过滤后的文件列表传入后续的复制/处理活动

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 18:10:03