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

如何在Dagster中实现多CSV并行、单CSV步骤串行的ETL工作流?

你的Dagster ETL场景完整实现方案

实现思路

  • 用get_csv_filenames从数据库拉取待处理的CSV文件列表
  • 靠generate_subtasks生成动态任务,实现不同CSV文件的并行处理
  • 给单个CSV定义串行处理流程:加载到DuckDB → 转换日期列 → 数字编码转文本分类 → 导出Parquet → 生成数据报告

完整代码实现

from typing import List
import duckdb
import pandas as pd
from dagster import (
    op,
    job,
    DynamicOut,
    DynamicOutput,
    In,
    Out,
    graph,
    Context
)

# 1. 查询数据库获取CSV文件列表
@op
def get_csv_filenames(context: Context) -> List[str]:
    # 替换成你的实际数据库查询逻辑,比如从PostgreSQL/SQLite拿CSV路径
    # 下面是模拟返回示例
    context.log.info("从数据库查询待处理CSV列表")
    return ["data/file1.csv", "data/file2.csv", "data/file3.csv"]

# 2. 生成动态子任务,每个CSV对应一个并行子流程
@op(out=DynamicOut(str))
def generate_subtasks(context: Context, csv_list: List[str]):
    for csv_filename in csv_list:
        # 把文件名里的特殊字符替换掉,作为mapping_key(Dagster要求唯一)
        mapping_key = csv_filename.replace("/", "_").replace(".", "_")
        context.log.info(f"为文件 {csv_filename} 创建并行子任务")
        yield DynamicOutput(csv_filename, mapping_key=mapping_key)

# 单个CSV的串行处理子图
@graph
def process_single_csv(csv_filename: str):
    # 加载CSV到DuckDB,返回表名供后续步骤使用
    table_name = load_csv_into_duckdb(csv_filename)
    # 转换日期列格式
    transformed_table = transform_dates(table_name)
    # 把数字编码转为文本分类
    categorized_table = from_code_2_categories(transformed_table)
    # 导出处理后的数据为Parquet
    parquet_path = export_2_parquet(categorized_table, csv_filename)
    # 生成该文件的数据概览报告
    generate_data_report(parquet_path, csv_filename)

# 子图里的各个操作步骤(OP)

@op(out=Out(str))
def load_csv_into_duckdb(context: Context, csv_filename: str) -> str:
    context.log.info(f"加载CSV文件 {csv_filename} 到DuckDB")
    # 连接DuckDB,这里用内存模式,也可以指定持久化文件路径
    conn = duckdb.connect()
    # 用文件名生成唯一表名,避免冲突
    table_name = f"csv_{csv_filename.split('/')[-1].replace('.csv', '')}"
    # 读取CSV到DuckDB表
    conn.execute(f"CREATE TABLE {table_name} AS SELECT * FROM read_csv('{csv_filename}')")
    conn.close()
    return table_name

@op(out=Out(str))
def transform_dates(context: Context, table_name: str) -> str:
    context.log.info(f"处理表 {table_name} 的日期列")
    conn = duckdb.connect()
    # 示例:假设日期列叫`date_str`,格式是'%Y-%m-%d',转为DATE类型
    # 你可以根据实际日期格式和列名修改
    conn.execute(f"""
        ALTER TABLE {table_name} 
        ALTER COLUMN date_str TYPE DATE 
        USING STRPTIME(date_str, '%Y-%m-%d')
    """)
    conn.close()
    return table_name

@op(out=Out(str))
def from_code_2_categories(context: Context, table_name: str) -> str:
    context.log.info(f"把表 {table_name} 的数字编码转为文本分类")
    conn = duckdb.connect()
    # 示例:假设`status_code`是数字编码列,映射为对应的文本分类
    # 你可以根据实际编码规则修改
    conn.execute(f"""
        ALTER TABLE {table_name}
        ADD COLUMN status_category VARCHAR;
        UPDATE {table_name}
        SET status_category = CASE
            WHEN status_code = 1 THEN '正常'
            WHEN status_code = 2 THEN '警告'
            WHEN status_code = 3 THEN '异常'
            ELSE '未知'
        END
    """)
    conn.close()
    return table_name

@op(out=Out(str))
def export_2_parquet(context: Context, table_name: str, csv_filename: str) -> str:
    context.log.info(f"将表 {table_name} 导出为Parquet文件")
    # 生成Parquet路径,和原CSV同目录,替换后缀
    parquet_path = csv_filename.replace(".csv", ".parquet")
    conn = duckdb.connect()
    # 导出表到Parquet
    conn.execute(f"COPY {table_name} TO '{parquet_path}' (FORMAT PARQUET)")
    conn.close()
    return parquet_path

@op
def generate_data_report(context: Context, parquet_path: str, csv_filename: str):
    context.log.info(f"生成 {csv_filename} 的数据概览报告")
    conn = duckdb.connect()
    # 读取Parquet文件到DataFrame
    df = conn.execute(f"SELECT * FROM read_parquet('{parquet_path}')").fetch_df()
    # 生成报告内容,包括行数、列数、缺失值统计等
    report_content = f"""
    数据概览报告 - {csv_filename}
    --------------------------
    总行数: {len(df)}
    列数: {len(df.columns)}
    缺失值统计:
    {df.isnull().sum().to_string()}
    """
    # 保存报告到文件
    report_path = csv_filename.replace(".csv", "_report.txt")
    with open(report_path, "w", encoding="utf-8") as f:
        f.write(report_content)
    context.log.info(f"报告已保存到 {report_path}")
    conn.close()

# 主Job:把所有流程串起来
@job
def etl_pipeline():
    csv_list = get_csv_filenames()
    # 把每个CSV分配到对应的串行子流程,实现并行处理
    generate_subtasks(csv_list).map(process_single_csv)

重点说明

  • 并行处理:generate_subtasks输出的DynamicOutput会让Dagster自动为每个CSV启动独立的并行子流程,能最大化利用机器资源
  • 串行执行保证:process_single_csv子图里的OP是按顺序执行的,前一步的输出作为后一步的输入,确保单个CSV的所有步骤严格串行
  • 灵活性:每个OP的逻辑都可以根据你的实际需求修改,比如数据库查询方式、日期格式、编码映射规则、报告内容等
  • 可调试性:每个OP都通过context.log输出关键日志,方便你排查处理过程中的问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 09:42:07