如何在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
相关产品推荐
相关产品推荐

