Dagster:I/O 管理器

Dagster 的 I/O 管理器可以把数据处理代码与读写数据的代码分开,减少重复代码,也方便更换数据存储位置。

很多 Dagster 流水线中的资产处理可以分为三步:

  1. 从数据存储读取到内存。
  2. 在内存中转换数据。
  3. 把转换结果写入数据存储。

对于遵循这种模式的资产,I/O 管理器可以简化从数据源读取和向其写入的代码。

开始之前

使用 Dagster 并不强制使用 I/O 管理器,它也不是所有场景的最佳选择。 如果每个资产的开头和结尾都重复加载、存储数据的相同代码,I/O 管理器可能有帮助,例如:

  • 资产保存在同一位置,并遵循一致的存储路径规则。
  • 本地、预发布和生产环境使用不同的存储方式。
  • 资产需要将上游依赖加载到内存后进行计算。

以下情况可能不适合使用 I/O 管理器:

  • 通过 SQL 查询直接在数据库中创建或更新表。
  • 流水线通过其他会写入存储的库或工具自行管理 I/O。
  • 资产无法装入内存,例如包含数十亿行的数据库表。

通常,如果为了使用 I/O 管理器反而让流水线更复杂,说明它可能不合适。这时应使用 deps 定义依赖关系。

在资产中使用 I/O 管理器

可以通过命令 dg scaffold defs dagster.asset <path/to/asset_file.py> 生成资产骨架。具体用法参见 dg CLI 文档。

下面的资产创建 DuckDB 连接,读取上游表,在内存中转换,再把结果写入 DuckDB 新表。

src/<project_name>/defs/assets.py:

import pandas as pd
from dagster_duckdb import DuckDBResource

import dagster as dg

raw_sales_data = dg.AssetSpec("raw_sales_data")

@dg.asset
def raw_sales_data(duckdb: DuckDBResource) -> None:
    # Read data from a CSV
    raw_df = pd.read_csv("https://docs.dagster.io/assets/raw_sales_data.csv")
    # Construct DuckDB connection
    with duckdb.get_connection() as conn:
        # Use the data from the CSV to create or update a table
        conn.execute(
            "CREATE TABLE IF NOT EXISTS raw_sales_data AS SELECT * FROM raw_df"
        )
        if not conn.fetchall():
            conn.execute("INSERT INTO raw_sales_data SELECT * FROM raw_df")
# Asset dependent on `raw_sales_data` asset
@dg.asset(deps=[raw_sales_data])
def clean_sales_data(duckdb: DuckDBResource) -> None:
    # Construct DuckDB connection
    with duckdb.get_connection() as conn:
        # Select data from table
        df = conn.execute("SELECT * FROM raw_sales_data").fetch_df()

        # Apply transform
        clean_df = df.fillna({"amount": 0.0})
        # Use transformed result to create or update a table
        conn.execute(
            "CREATE TABLE IF NOT EXISTS clean_sales_data AS SELECT * FROM clean_df"
        )
        if not conn.fetchall():
            conn.execute("INSERT INTO clean_sales_data SELECT * FROM clean_df")
# Asset dependent on `clean_sales_data` asset
@dg.asset(deps=[clean_sales_data])
def sales_summary(duckdb: DuckDBResource) -> None:
    # Construct DuckDB connection
    with duckdb.get_connection() as conn:
        # Select data from table
        df = conn.execute("SELECT * FROM clean_sales_data").fetch_df()

        # Apply transform
        summary = df.groupby(["owner"])["amount"].sum().reset_index()
        # Use transformed result to create or update a table
        conn.execute(
            "CREATE TABLE IF NOT EXISTS sales_summary AS SELECT * from summary"
        )
        if not conn.fetchall():
            conn.execute("INSERT INTO sales_summary SELECT * from summary")

src/<project_name>/defs/resources.py:

from dagster_duckdb import DuckDBResource

import dagster as dg


@dg.definitions
def resources():
    return dg.Definitions(
        resources={"duckdb": DuckDBResource(database="sales.duckdb", schema="public")}
    )

使用 I/O 管理器后,可以从资产中移除读写代码,把这部分工作交给管理器。资产只保留转换逻辑或获取初始 CSV 的代码。

src/<project_name>/defs/assets.py:

import pandas as pd
from dagster_duckdb_pandas import DuckDBPandasIOManager

import dagster as dg


@dg.asset
def raw_sales_data() -> pd.DataFrame:
    return pd.read_csv("https://docs.dagster.io/assets/raw_sales_data.csv")
# highlight-start
@dg.asset
# Load the upstream `raw_sales_data` asset as input & specify the returned data type (`pd.DataFrame`)
def clean_sales_data(raw_sales_data: pd.DataFrame) -> pd.DataFrame:
    # Storing data with an I/O manager requires returning the data
    return raw_sales_data.fillna({"amount": 0.0})
    # highlight-end


@dg.asset
def sales_summary(clean_sales_data: pd.DataFrame) -> pd.DataFrame:
    return clean_sales_data.groupby(["owner"])["amount"].sum().reset_index()

src/<project_name>/defs/resources.py:

from dagster_duckdb_pandas import DuckDBPandasIOManager

import dagster as dg


@dg.definitions
def resources():
    return dg.Definitions(
        # highlight-start
        # Define the I/O manager and pass it to `Definitions`
        resources={
            "io_manager": DuckDBPandasIOManager(
                database="sales.duckdb", schema="public"
            )
        }
        # highlight-end
    )

要通过 I/O 管理器加载上游资产,把它声明为资产函数的输入参数。本例中,DuckDBPandasIOManager 读取与上游资产同名的 DuckDB 表 raw_sales_data,并把数据作为 Pandas DataFrame 传给 clean_sales_data。

要通过 I/O 管理器存储数据,在资产函数中返回数据即可。返回值必须是该管理器支持的类型。本例返回 Pandas DataFrame,管理器将其写入与资产同名的 DuckDB 表。

各种管理器支持哪些类型、如何保存数据,参见各自文档。

更换数据存储

有了 I/O 管理器,更换存储只需替换管理器实现。只包含转换逻辑的资产定义不必修改。

下面使用 Snowflake I/O 管理器替换 DuckDB 管理器。

src/<project_name>/defs/assets.py:

import pandas as pd

import dagster as dg


@dg.asset
def raw_sales_data() -> pd.DataFrame:
    return pd.read_csv("https://docs.dagster.io/assets/raw_sales_data.csv")


@dg.asset
def clean_sales_data(raw_sales_data: pd.DataFrame) -> pd.DataFrame:
    return raw_sales_data.fillna({"amount": 0.0})


@dg.asset
def sales_summary(clean_sales_data: pd.DataFrame) -> pd.DataFrame:
    return clean_sales_data.groupby(["owner"])["amount"].sum().reset_index()

src/<project_name>/defs/resources.py:

from dagster_snowflake_pandas import SnowflakePandasIOManager

import dagster as dg

@dg.definitions
def resources():
    return dg.Definitions(
        resources={
            # highlight-start
            # Swap in a Snowflake I/O manager
            "io_manager": SnowflakePandasIOManager(
                database=dg.EnvVar("SNOWFLAKE_DATABASE"),
                account=dg.EnvVar("SNOWFLAKE_ACCOUNT"),
                user=dg.EnvVar("SNOWFLAKE_USER"),
                password=dg.EnvVar("SNOWFLAKE_PASSWORD"),
            )
            # highlight-end
        }
    )

内置 I/O 管理器

Dagster 为常见数据存储和内存数据格式提供了现成的 I/O 管理器实现。

名称 说明
FilesystemIOManager 默认 I/O 管理器,把输出作为 pickle 文件保存在本地文件系统
InMemoryIOManager 把输出保存在内存,主要用于单元测试
s3.S3PickleIOManager 把输出作为 pickle 文件保存在 AWS S3
adls2.ConfigurablePickledObjectADLS2IOManager 把输出作为 pickle 文件保存在 Azure ADLS2
GCSPickleIOManager 把输出作为 pickle 文件保存在 Google Cloud Platform GCS
BigQueryPandasIOManager 把 Pandas DataFrame 输出保存在 Google Cloud Platform BigQuery
BigQueryPySparkIOManager 把 PySpark DataFrame 输出保存在 Google Cloud Platform BigQuery
SnowflakePandasIOManager 把 Pandas DataFrame 输出保存在 Snowflake
SnowflakePySparkIOManager 把 PySpark DataFrame 输出保存在 Snowflake
ClickhousePandasIOManager 把 Pandas DataFrame 输出保存在 ClickHouse
ClickhousePolarsIOManager 把 Polars DataFrame 输出保存在 ClickHouse
DuckDBPandasIOManager 把 Pandas DataFrame 输出保存在 DuckDB
DuckDBPySparkIOManager 把 PySpark DataFrame 输出保存在 DuckDB
DuckDBPolarsIOManager 把 Polars DataFrame 输出保存在 DuckDB

后续步骤


来源:I/O managers,Dagster 官方文档,Dagster Labs 与贡献者。本文为中文翻译,完整展开原文引用的6个Python文件,保留代码及其中的高亮标记注释。

Copyright 2025 Dagster Labs, Inc.;原文网站页脚 Copyright 2026 Dagster Labs。采用 Apache License 2.0。除非法律要求或书面约定,按该许可发布的作品均按原样提供,不附带任何明示或默示保证或条件。具体权利与限制参见许可证全文。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容