Python脚本如何生成Dagster配置

wen 实用脚本 22

本文目录导读:

Python脚本如何生成Dagster配置

  1. 使用Python字典直接定义配置
  2. 动态生成配置JSON
  3. 使用ConfigMapping动态生成
  4. 从外部源生成配置
  5. 运行时的配置生成示例
  6. 配置验证和转换
  7. 最佳实践建议

我来详细介绍如何生成Dagster配置的几种方法:

使用Python字典直接定义配置

from dagster import op, job, Config
# 定义配置类
class MyOpConfig(Config):
    batch_size: int = 100
    learning_rate: float = 0.01
    model_path: str = "/models/model.pkl"
@op
def my_operation(config: MyOpConfig):
    # 使用配置
    print(f"Batch size: {config.batch_size}")
    print(f"Learning rate: {config.learning_rate}")
@job
def my_job():
    my_operation()

动态生成配置JSON

import json
from dagster import RunConfig, ConfigMapping
from typing import Dict, Any
class ConfigGenerator:
    @staticmethod
    def generate_training_config(model_name: str, params: Dict) -> Dict:
        """动态生成训练配置"""
        config = {
            "ops": {
                "process_data": {
                    "config": {
                        "input_path": f"data/{model_name}/input.csv",
                        "batch_size": params.get("batch_size", 32)
                    }
                },
                "train_model": {
                    "config": {
                        "model_name": model_name,
                        "learning_rate": params.get("learning_rate", 0.001),
                        "epochs": params.get("epochs", 10)
                    }
                }
            }
        }
        return config
    @staticmethod
    def save_config_to_file(config: Dict, filepath: str):
        """保存配置到JSON文件"""
        with open(filepath, 'w') as f:
            json.dump(config, f, indent=2)

使用ConfigMapping动态生成

from dagster import ConfigMapping
def generate_config_resolver():
    """创建配置映射函数"""
    def my_config_mapping(config: Dict) -> Dict:
        """根据输入参数动态生成配置"""
        # 解析输入参数
        env = config.get("environment", "dev")
        data_cfg = config.get("data", {})
        model_cfg = config.get("model", {})
        # 生成具体配置
        generated_config = {
            "ops": {
                "load_data": {
                    "config": {
                        "source": f"s3://{env}-bucket/data/{data_cfg.get('dataset', 'default')}",
                        "format": data_cfg.get("format", "csv")
                    }
                },
                "train_model": {
                    "config": {
                        "model_type": model_cfg.get("type", "xgboost"),
                        "hyperparameters": {
                            "learning_rate": model_cfg.get("lr", 0.1),
                            "max_depth": model_cfg.get("max_depth", 6)
                        }
                    }
                }
            },
            "resources": {
                "io_manager": {
                    "config": {
                        "base_path": f"/data/{env}/results/"
                    }
                }
            }
        }
        return generated_config
    return ConfigMapping(my_config_mapping)

从外部源生成配置

import yaml
import os
from dagster import RunConfig
class ExternalConfigGenerator:
    """从外部源生成Dagster配置"""
    @staticmethod
    def from_yaml_file(yaml_path: str) -> Dict:
        """从YAML文件读取配置"""
        with open(yaml_path, 'r') as f:
            config = yaml.safe_load(f)
        return config
    @staticmethod
    def from_env_variables(prefix: str = "DAGSTER_") -> Dict:
        """从环境变量生成配置"""
        config = {
            "resources": {
                "db": {
                    "config": {
                        "host": os.getenv(f"{prefix}DB_HOST", "localhost"),
                        "port": int(os.getenv(f"{prefix}DB_PORT", "5432")),
                        "database": os.getenv(f"{prefix}DB_NAME", "dagster")
                    }
                }
            }
        }
        return config
    @staticmethod
    def from_template(template_path: str, variables: Dict) -> Dict:
        """从模板生成配置"""
        import string
        with open(template_path, 'r') as f:
            template_content = f.read()
        # 填充变量
        template = string.Template(template_content)
        filled_content = template.safe_substitute(variables)
        # 解析为字典
        config = json.loads(filled_content)
        return config

运行时的配置生成示例

from dagster import job, op, RunConfig, DagsterInstance
from datetime import datetime
@op(required_resource_keys={"config_generator"})
def generate_dynamic_config(context):
    """在运行时动态生成配置"""
    # 获取运行时信息
    run_id = context.run_id
    current_time = datetime.now()
    # 生成配置
    config = {
        "ops": {
            "process_data": {
                "config": {
                    "timestamp": current_time.isoformat(),
                    "run_id": run_id,
                    "data_version": f"v{current_time.strftime('%Y%m%d_%H%M%S')}"
                }
            }
        }
    }
    return config
# 使用配置启动作业
def launch_job_with_dynamic_config():
    instance = DagsterInstance.get()
    # 生成配置
    config = {
        "ops": {
            "sample_op": {
                "config": {
                    "param1": "value1",
                    "param2": 100
                }
            }
        }
    }
    # 启动作业
    result = my_job.execute_in_process(
        run_config=RunConfig(**config)
    )
    return result

配置验证和转换

from pydantic import BaseModel, validator
from typing import Optional
class DataProcessingConfig(BaseModel):
    """数据处理的配置模式"""
    input_path: str
    output_path: str
    batch_size: int = 100
    parallel: bool = True
    @validator('batch_size')
    def validate_batch_size(cls, v):
        if v <= 0:
            raise ValueError("Batch size must be positive")
        if v > 10000:
            raise ValueError("Batch size too large")
        return v
def generate_validated_config(base_config: Dict) -> Dict:
    """生成并验证配置"""
    # 验证配置
    validated_config = DataProcessingConfig(**base_config)
    # 生成完整的Dagster配置
    dagster_config = {
        "ops": {
            "data_processing": {
                "config": validated_config.dict()
            }
        }
    }
    return dagster_config
# 使用示例
config_data = {
    "input_path": "/data/input.csv",
    "output_path": "/data/output.parquet",
    "batch_size": 500
}
validated_config = generate_validated_config(config_data)

最佳实践建议

  1. 使用类型提示:利用Python的类型提示来定义配置结构
  2. 配置验证:使用Pydantic等库进行配置验证
  3. 分层配置:区分开发、测试和生产环境的配置
  4. 配置版本控制:将关键配置纳入版本控制
  5. 可序列化:确保配置可以安全地序列化为JSON/YAML

这些方法可以根据你的具体需求进行组合和扩展,以生成符合Dagster要求的配置。

抱歉,评论功能暂时关闭!