Dagster Pipeline Analyzer — Dagster 流水线分析器
v1.0.0分析Dagster管道和软件定义资产的质量、调度、分区、IO管理器、资源配置和可观察性。检查一个...
运行时依赖
安装命令
点击复制本土化适配说明
Dagster Pipeline Analyzer — Dagster 流水线分析器 安装说明: 安装命令:["openclaw skills install dagster-pipeline-analyzer"] 支持国内镜像加速,使用 --registry https://cn.longxiaskill.com 参数可加速下载
技能文档
Dagster Pipeline Analyzer 分析Dagster软件定义的资产和管道的质量、可靠性和运营卓越性。 审查资产依赖图、新鲜度策略、分区策略、IO管理器、资源、传感器/调度和可观察性差距。 作为一名高级数据平台工程师,审计您的Dagster部署。
使用 基本:分析位于/path/to/dagster/的Dagster项目 聚焦:检查资产新鲜度策略 | 分析分区策略 | 审查IO管理器配置 | 查找没有可观察性元数据的资产
工作原理 步骤1:发现Dagster定义 find /path/to/project -name "definitions.py" -o -name "repository.py" find /path/to/project -name ".py" -path "/assets/" find /path/to/project -name ".py" -path "/sensors/" -o -path "/schedules/" cat /path/to/project/dagster.yaml /path/to/project/workspace.yaml 2>/dev/null 解析@asset、@multi_asset、@op、@graph、@job、ConfigurableResource、ConfigurableIOManager、@schedule、@sensor和分区定义。
步骤2:审计资产图结构 资产图:34个资产,5个组(raw、staging、warehouse、analytics、ml) 最大深度:7 | 外部资产:3 | 源资产:4 [source] s3_raw_events -> raw_events -> cleaned_events -> event_aggregates -> daily_metrics -> churn_features -> churn_model 失败:"orphan_transform"没有下游消费者 上次物化49天前。 删除或记录目的。 失败:"daily_metrics"依赖于6个上游资产 —— 脆弱瓶颈 如果任何上游失败,物化被阻塞。 推荐:添加重试逻辑或部分物化支持 失败:没有为任何资产定义@asset_check 推荐:添加行数、空检查、新鲜度、模式验证
步骤3:审查新鲜度策略 具有新鲜度策略的资产:8/34(24%) 失败:"daily_revenue" —— 没有业务关键指标的新鲜度策略 修复:FreshnessPolicy(maximum_lag_minutes=120) 失败:新鲜度SLA级联违规:"event_aggregates"具有60分钟新鲜度,但上游"cleaned_events"需要~20分钟才能物化 有效预算:10分钟。 实际平均:25分钟。 SLA经常被违反。 修复:放松到90分钟或加快物化 警告:"user_segments"具有1440分钟(24小时)新鲜度,但下游"real_time_recommendations"期望新鲜数据 修复:紧缩到maximum_lag_minutes=60
步骤4:分析分区 "raw_events":DailyPartitionsDefinition —— 通过,匹配数据模式 "event_aggregates":MonthlyPartitionsDefinition 警告:上游是每日的 —— 验证TimeWindowPartitionMapping "ml_features":未分区,处理2M+行每次运行(45分钟) 失败:添加DailyPartitionsDefinition进行增量处理(~3分钟/分区) 失败:没有BackfillPolicy在任何分区资产上 风险:意外回填486个分区 = 486个并发运行 修复:BackfillPolicy(max_partitions_per_run=10)
步骤5:审查IO管理器 配置:"io_manager"(Filesystem)、"warehouse_io"(Snowflake)、"s3_io"(S3Pickle) 失败:默认是FilesystemIOManager —— 数据在pod重启时丢失 修复:设置默认为S3/GCS/数据库支持的存储 失败:S3PickleIOManager —— 不可移植,安全风险(任意代码执行) 修复:切换到S3ParquetIOManager以执行模式强制和压缩 警告:没有返回类型注释 —— IO管理器无法验证模式 修复:@asset def daily_metrics(...) -> pd.DataFrame: 覆盖率:warehouse_io 35%、s3_io 24%、filesystem 41%(令人担忧)
步骤6:审计资源 失败:snowflake_password="..."在definitions.py第42行 修复:使用EnvVar("SNOWFLAKE_PASSWORD") 失败:资源"spark"定义但从未使用 —— 删除它 警告:资源"dbt"具有硬编码的项目路径 修复:使用EnvVar或相对路径 警告:没有资源级别的运行时健康检查
步骤7:审查调度和传感器 失败:"daily_etl"使用cron,但资产具有DailyPartitionsDefinition 修复:build_schedule_from_partitions()用于自动对齐 失败:"s3_file_sensor"每30秒轮询一次 —— 高API成本,速率限制风险 修复:minimum_interval_seconds=300或SQS/SNS事件驱动 失败:"freshness_sensor"没有错误处理 —— 崩溃传感器守护进程 修复:包装在try/except中,在失败时返回SkipReason 警告:14个资产没有自动化 —— 手动仅物化 警告:没有@asset_sensor用于跨作业依赖协调
步骤8:检查可观察性 失败:22/34个资产没有描述 失败:没有MaterializeResult元数据(没有行数、模式、质量) 修复:yield MaterializeResult(metadata={"row_count": len(df), ...}) 失败:没有为任何资产定义资产检查 警告:没有代码版本在任何资产上 —— 始终重新物化 警告:没有定义所有者 —— 警报不可路由 覆盖率:描述35%、新鲜度24%、代码版本0%、检查0%、元数据15%、所有者0%
步骤9:最终报告 # Dagster Pipeline Analysis Report