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

如何在Dagster中以Asset为Op输入并高效处理批量数据?

问题背景与疑问

我从Dagster文档了解到,可使用op和graph生成asset,但希望将asset作为op的输入,具体业务场景如下:

  • 从外部API获取国家元数据列表并存储为资源:
@dagster.asset
def country_metadata_asset() -> List[Dict]:
    ...
  • 基于该asset定义下游asset:
@dagster.asset
def country_names_asset(country_metadata_asset) -> List[str]:
    ...
  • 希望使用该asset调用另一数据源,检索并验证数据后写入资源,但返回数据量极大,需批量处理,因此考虑用graph结合op实现,尝试代码如下:
@dagster.op(out=dagster.DynamicOut())
def load_country_names(country_names_asset):
    for country_index, country_name in enumerate(country_names_asset):
        yield dagster.DynamicOutput(
            country_name, mapping_key=f"{country_index} {country_name}"
        )

@dagster.graph()
def update_data_graph():
    country_names = load_country_names()
    country_names.map(retrieve_and_process_data)

@dagster.job()
def run_update_job():
    update_data_graph()

但该方案无法运行,且不确定是否符合概念逻辑,现咨询两个问题:

  1. 如何让load_country_names的输入识别为asset?是否需要在op内手动物化?
  2. 如何高效将retrieve_and_process_data返回的增量数据写入资源?因无法将数据保存在内存中,考虑用自定义IOManager但不知具体实现方式。

问题解答

1. 让Op识别Asset输入的方法

不需要在Op内手动物化,核心是将Graph转换为Graph-backed Asset,同时通过AssetIn声明Op依赖的上游Asset:

修正后代码示例

from dagster import (
    op, graph, DynamicOut, DynamicOutput, asset, AssetIn, graph_asset
)
from typing import List, Dict

# 原有的asset定义保持不变
@asset
def country_metadata_asset() -> List[Dict]:
    ...

@asset
def country_names_asset(country_metadata_asset) -> List[str]:
    ...

@op(out=DynamicOut())
def load_country_names(country_names):
    for country_index, country_name in enumerate(country_names):
        yield DynamicOutput(
            country_name, mapping_key=f"{country_index} {country_name}"
        )

@op
def retrieve_and_process_data(country_name):
    # 实现单国家的数据检索与处理逻辑
    ...
    return processed_data

@graph
def update_data_graph(country_names):
    names = load_country_names(country_names)
    return names.map(retrieve_and_process_data)

# 将Graph转为Graph-backed Asset,明确关联上游的country_names_asset
@graph_asset(ins={"country_names": AssetIn("country_names_asset")})
def processed_country_data_asset():
    return update_data_graph()

关键说明

  • 用graph_asset替代普通graph,让整个Graph成为一个Asset,Dagster会自动处理上游依赖的加载与物化。
  • 通过AssetIn指定输入对应的Asset名称,Op会直接接收该Asset的物化结果作为输入,无需手动调用物化方法。

2. 自定义IOManager实现增量数据写入

针对大体积增量数据,自定义IOManager可实现边处理边写入,避免内存溢出,核心是在handle_output方法中实现增量写入逻辑:

自定义IOManager示例

from dagster import IOManager, OutputContext, InputContext
import pandas as pd
import os
from typing import Iterator

class IncrementalFileIOManager(IOManager):
    def __init__(self, base_path: str):
        self.base_path = base_path
        os.makedirs(base_path, exist_ok=True)

    def handle_output(self, context: OutputContext, obj: Iterator):
        # 基于动态输出的mapping_key生成独立文件,避免单文件过大
        mapping_key = context.mapping_key
        file_path = f"{self.base_path}/{mapping_key}.parquet"
        
        # 逐批写入增量数据
        for batch in obj:
            if not os.path.exists(file_path):
                batch.to_parquet(file_path, index=False)
            else:
                # 追加模式写入已有文件
                existing_data = pd.read_parquet(file_path)
                pd.concat([existing_data, batch]).to_parquet(file_path, index=False)

    def load_input(self, context: InputContext):
        # 加载时读取对应mapping_key的文件
        file_path = f"{self.base_path}/{context.mapping_key}.parquet"
        return pd.read_parquet(file_path)

配置与使用

  1. 在Definitions中注册IOManager:
from dagster import Definitions

defs = Definitions(
    assets=[country_metadata_asset, country_names_asset, processed_country_data_asset],
    resources={"io_manager": IncrementalFileIOManager(base_path="./processed_country_data")}
)
  1. 修改retrieve_and_process_data Op,返回增量迭代器:
@op
def retrieve_and_process_data(country_name) -> Iterator[pd.DataFrame]:
    # 模拟按批次获取并处理数据,返回迭代器
    for data_batch in fetch_country_data_in_batches(country_name):
        processed_batch = clean_validate_data(data_batch)
        yield processed_batch

关键说明

  • IOManager的handle_output接收Op返回的迭代器,逐批写入存储,避免一次性加载所有数据到内存。
  • 利用动态输出的mapping_key为每个国家生成独立文件,支持并行处理,也便于后续单独查询某国数据。
  • 可根据实际存储介质(如S3、MySQL)调整写入逻辑,比如数据库可使用批量插入语句替代文件写入。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 05:30:53