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

如何从Dagster RunFailureSensorContext获取run_config处理失败文件?

问题:在Dagster的run_failure_sensor中获取失败任务的run_config并处理失败文件

我已经在Dagster中配置了触发excel_to_csv任务的sensor:

@sensor(job=excel_to_csv)
def convert_unit_list(ctx: SensorEvaluationContext):

    env = os.getenv("ENVIRONMENT")
    region = os.getenv("AWS_REGION")
    appcfg = AppConfig("unit-ingestion", "developer-units", env, region)

    bucket = Bucket.LANDING
    bucket_name = bucket.format(env)
    prefix = "unit-ingestion/developer-units"
    ctx.log.info(f"Searching {bucket_name}/{prefix} for new files...")

    since_key = ctx.cursor
    keys = get_s3_keys(bucket_name, prefix, since_key)
    if not keys:
        return SkipReason(f"No keys found in {bucket_name}/{prefix}")

    run_config = {
        "s3": S3Resource(region_name=region),
        "conv_ctx": ExcelConversionContext(appcfg, bucket, "unit-ingestion/units", 1, "B:BV"),
    }

    ctx.log.info(f"Found {len(keys)} new keys in {bucket_name}/{prefix}")
    for key in keys:

        yield RunRequest(key, run_config=run_config)

        ctx.log.info(f"Updating cursor to {key}")
        ctx.update_cursor(key)

现在需要编写run_failure_sensor,在excel_to_csv任务失败时将对应的Excel文件移动到S3拒绝存储桶。目前的基础代码如下,但不知道如何获取失败任务的run_config:

@run_failure_sensor(request_job=excel_to_csv)
def excel_to_csv_failure(ctx: RunFailureSensorContext):

    ctx.log.error(f"{ctx.sensor_name} failed to process {ctx.failure_event.asset_key}")

解决方案

核心要点

  1. 不要依赖ctx.dagster_event.step_input_data:这个字段是步骤级别的输入数据,并非全局的run_config。正确获取run_config的方式是通过Dagster实例获取完整的运行记录。
  2. 利用RunRequest的run_key:之前的sensor中,每个RunRequest的key参数会作为运行的dagster/run_key标签存储,通过这个标签可以直接拿到失败的S3文件路径,比从run_config中提取更直接。

修改后的完整代码

import os
from dagster import run_failure_sensor, RunFailureSensorContext
from your_module import S3Resource, Bucket  # 替换为你的实际模块路径

@run_failure_sensor(request_job=excel_to_csv)
def excel_to_csv_failure(ctx: RunFailureSensorContext):
    # 获取失败任务的运行记录
    run = ctx.instance.get_run_by_id(ctx.dagster_event.run_id)
    
    # 从run的标签中拿到对应的S3文件key(即之前sensor中yield RunRequest时传入的key)
    failed_file_key = run.tags.get("dagster/run_key")
    if not failed_file_key:
        ctx.log.error("无法从运行标签中获取失败文件的S3路径")
        return
    
    # 获取源存储桶和目标拒绝存储桶
    env = os.getenv("ENVIRONMENT")
    source_bucket = Bucket.LANDING.format(env)
    reject_bucket = Bucket.REJECT.format(env)  # 需确保你定义了REJECT类型的Bucket枚举
    
    # 初始化S3客户端
    region = os.getenv("AWS_REGION")
    s3_client = S3Resource(region_name=region).get_client()
    
    # 移动文件:先复制到拒绝桶,再删除源文件
    target_key = f"failed-units/{failed_file_key}"
    try:
        # 复制文件
        s3_client.copy_object(
            Bucket=reject_bucket,
            Key=target_key,
            CopySource={"Bucket": source_bucket, "Key": failed_file_key}
        )
        # 删除源文件
        s3_client.delete_object(Bucket=source_bucket, Key=failed_file_key)
        ctx.log.info(f"已将失败文件 {failed_file_key} 从 {source_bucket} 移动到 {reject_bucket}/{target_key}")
    except Exception as e:
        ctx.log.error(f"移动失败文件 {failed_file_key} 时出错: {str(e)}")
    
    ctx.log.error(f"{ctx.sensor_name} 处理文件 {failed_file_key} 失败")

补充说明

  • 如果确实需要获取run_config,可以直接通过run.run_config拿到完整的配置内容,比如提取conv_ctx中的参数:conv_ctx_config = run.run_config.get("conv_ctx", {})。
  • 移动文件使用S3的copy_object+delete_object组合,避免了本地下载再上传的额外开销,效率更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 12:32:36