如何在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
相关产品推荐
相关产品推荐

