Dagster如何将官方s3_resource传入自定义封装的HelperAwsS3资源类
你的实现思路是合理的,只需要给自定义resource声明依赖的官方S3资源即可,具体修改如下:
步骤1:修改自定义资源声明,指定依赖资源
给@resource装饰器添加required_resource_keys参数,声明当前资源依赖官方s3_resource,然后从context.resources中取出初始化完成的s3实例传入你封装的工具类即可:
@resource(required_resource_keys={"s3"}) def connection_helper_aws_s3_resource(context): # context.resources.s3 就是dagster_aws提供的已完成鉴权的s3资源实例 return HelperAwsS3(s3_resource=context.resources.s3)
步骤2:模式定义中注册两个资源
在你的ModeDefinition里同时注册官方s3资源和你封装的helper资源:
run_mode = ModeDefinition( resource_defs={ "s3": s3_resource, # 官方S3资源 "s3_helper": connection_helper_aws_s3_resource # 你封装的helper资源 } )
步骤3:在Solid中调用封装好的方法
给需要用到S3操作的solid声明依赖s3_helper资源,即可直接调用你封装的各类方法:
@solid(required_resource_keys={"s3_helper"}) def list_s3_files_solid(context, bucket_name: str, prefix: str): res = context.resources.s3_helper.s3_list_bucket(bucket_name, prefix) context.log.info(f"桶内匹配文件数:{len(res.get('Contents', []))}") # 其余上传、下载方法调用逻辑同理
补充说明
你当前的封装思路非常适合跨Pipeline复用通用S3操作,没有必要更换实现方式。如果后续需要对接其他存储(比如OSS、本地文件系统),也可以用相同的思路封装统一的操作接口,上层Solid不需要修改任何逻辑就能切换存储源。
内容的提问来源于stack exchange,提问作者Gayathri
相关产品推荐
相关产品推荐

