Dagster技术问询:如何通过外部API获取仅修改资产并传递日期参数
在Dagster中存储并传递上次运行日期的方案
针对你需要查询API增量数据的需求,这里有三种落地性强的方案,结合Dagster的特性实现:
方案1:利用Dagster内置的Asset元数据存储
这是最轻量化的方式,无需额外依赖,直接用Dagster自带的Asset元数据记录上次查询时间:
- 每次Asset运行完成后,把本次查询到的最大修改时间(或运行结束时间)写入Asset的元数据
- 下次运行时,从Dagster实例中读取该Asset的最新物化记录,提取日期值作为API查询的起始点
示例代码:
from dagster import asset, AssetExecutionContext from my_api_resource import my_api_client # 你定义的API调用resource @asset def incremental_api_data(context: AssetExecutionContext, api_client: my_api_client): # 读取上次运行的时间 last_run_time = None latest_materialization = context.instance.get_latest_materialization(context.asset_key) if latest_materialization: last_run_time = latest_materialization.metadata.get("last_modified_time").value # 调用API,传入过滤条件 data = api_client.fetch_data(since=last_run_time) # 处理并保存数据 process_and_save_data(data) # 记录本次的最大修改时间到元数据 if data: max_modified_time = max(item["modified_at"] for item in data) context.add_output_metadata({"last_modified_time": max_modified_time})
方案2:自定义状态跟踪Resource
如果需要更可靠的持久化存储(比如跨Dagster实例共享状态、或状态不能丢失),可以写一个专门的Resource来读写运行时间:
- 该Resource可对接本地文件、SQL数据库或KV存储(如Redis)
- Asset中依赖该Resource,运行时先读取历史时间,查询完成后更新状态
示例代码:
from dagster import resource, asset, AssetExecutionContext import sqlite3 @resource def last_modified_tracker(): class Tracker: def __init__(self): self.conn = sqlite3.connect("tracker.db") self.conn.execute("CREATE TABLE IF NOT EXISTS run_times (asset_key TEXT PRIMARY KEY, last_time TEXT)") def read(self, asset_key): cursor = self.conn.execute("SELECT last_time FROM run_times WHERE asset_key = ?", (asset_key,)) result = cursor.fetchone() return result[0] if result else None def write(self, asset_key, last_time): self.conn.execute("REPLACE INTO run_times (asset_key, last_time) VALUES (?, ?)", (asset_key, last_time)) self.conn.commit() return Tracker() @asset def incremental_api_data(context: AssetExecutionContext, api_client: my_api_client, tracker: last_modified_tracker): asset_key_str = ".".join(context.asset_key.path) last_run_time = tracker.read(asset_key_str) data = api_client.fetch_data(since=last_run_time) process_and_save_data(data) if data: max_modified_time = max(item["modified_at"] for item in data) tracker.write(asset_key_str, max_modified_time)
方案3:使用分区化Asset(Partitioned Assets)
如果你的数据本身按时间自然分区(比如按天/小时生成),直接用Dagster的分区功能:
- 定义按时间分区的Asset,每个分区对应一个时间窗口
- 用当前分区的起始时间作为API查询的过滤条件,或直接用分区键作为时间参数
- 这种方式天生自带时间跟踪,无需额外存储状态
示例代码:
from dagster import asset, DailyPartitionsDefinition from my_api_resource import my_api_client daily_partitions = DailyPartitionsDefinition(start_date="2024-01-01") @asset(partitions_def=daily_partitions) def daily_api_data(context: AssetExecutionContext, api_client: my_api_client): # 获取当前分区的日期 partition_date = context.partition_key # 按日期区间查询API数据 start_time = f"{partition_date}T00:00:00" end_time = f"{partition_date}T23:59:59" data = api_client.fetch_data(since=start_time, until=end_time) process_and_save_data(data, partition_date)
方案选择建议
- 优先选方案1:轻量、无需额外运维,适合大多数增量同步场景
- 方案2:适合跨实例共享状态、或状态需要持久化不丢失的场景
- 方案3:适合数据本身有明确时间分区规则的ETL流程
另外,关于API调用的实现:推荐封装成Resource,这样可以统一管理API配置(如密钥、Base URL),方便在多个Asset中复用,也更容易测试和维护。
内容的提问来源于stack exchange,提问作者beginner_
相关产品推荐
相关产品推荐

