如何从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}")
解决方案
核心要点
- 不要依赖
ctx.dagster_event.step_input_data:这个字段是步骤级别的输入数据,并非全局的run_config。正确获取run_config的方式是通过Dagster实例获取完整的运行记录。 - 利用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
相关产品推荐
相关产品推荐

