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

如何在Dagster中基于记录ID实现增量加载?

在Dagster中实现HTTP API的增量加载

要实现基于fromUpdateId参数的增量加载,你可以通过资产物化元数据记录每次获取的最高更新ID,具体实现方案如下:

核心逻辑

每次执行资产时,先从历史物化元数据中读取上次记录的最高updateId,以此作为fromUpdateId参数调用API;获取数据后提取本次的最高ID,将其写入本次物化的元数据,供下次执行时复用。

具体代码实现

from dagster import asset, AssetExecutionContext
import requests

@asset
def incremental_change_log(context: AssetExecutionContext):
    # 1. 读取历史元数据中的最高updateId,首次执行用默认值0(可根据API规则调整)
    last_update_id = context.metadata.get("highest_update_id", 0)
    
    # 2. 调用目标API,传入fromUpdateId参数
    api_response = requests.get(
        "https://your-api-domain.com/changes",
        params={"fromUpdateId": last_update_id}
    )
    api_response.raise_for_status()
    change_data = api_response.json()
    
    # 3. 提取本次返回数据中的最高updateId
    if not change_data:
        highest_update_id = last_update_id  # 无新数据时保持原ID,避免下次重复请求
    else:
        highest_update_id = max(item["updateId"] for item in change_data)
    
    # 4. 自定义数据处理逻辑(如写入数据库、存储至数据湖等,此处省略)
    # process_and_save_data(change_data)
    
    # 5. 将最高updateId写入本次物化元数据
    return {
        "metadata": {
            "highest_update_id": highest_update_id,
            "new_records_count": len(change_data)
        }
    }

关键细节说明

  • 历史元数据读取:context.metadata会自动加载该资产上次物化时的元数据,首次执行时因无历史记录,会使用你设置的默认值。
  • 空数据处理:当API返回空列表(无新更新)时,保持highest_update_id不变,避免下次执行传入无效参数。
  • 元数据持久化:Dagster会自动维护资产的物化元数据,无需额外搭建存储服务,下次执行可直接读取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 16:05:48