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

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_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 17:50:23