Dagster 的 I/O 管理器可以把数据处理代码与读写数据的代码分开,减少重复代码,也方便更换数据存储位置。
很多 Dagster 流水线中的资产处理可以分为三步:
- 从数据存储读取到内存。
- 在内存中转换数据。
- 把转换结果写入数据存储。
对于遵循这种模式的资产,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。除非法律要求或书面约定,按该许可发布的作品均按原样提供,不附带任何明示或默示保证或条件。具体权利与限制参见许可证全文。











暂无评论内容