OSS使用

更新时间:
复制 MD 格式

PAIDLCDSW支持通过ossfs 2.0JindoFuse挂载OSS到容器路径,也支持通过OSS Connector for AI/MLOSS SDK直接读取OSS数据。按是否使用PyTorch、是否需要POSIX语义等场景选择合适方式。

背景信息

AI开发常把源数据存到OSS,再下载到训练环境。这会带来以下问题:

  • 数据集下载时间长,GPU空等。

  • 每次训练任务都需要重复下载数据。

  • 随机采样要求每个训练节点都下载完整数据集。

参考下表选择合适的OSS读取方式:

OSS数据读取方法

描述

适用场景

JindoFuse

通过JindoFuseOSS数据集挂载到容器路径,直接读写数据。

  • 希望像访问本地数据一样读取OSS,或数据集较小、可用JindoFuse本地缓存加速。

  • 使用的框架不是PyTorch。

  • 需要向OSS写入数据。

ossfs 2.0

ossfs 2.0是一款高性能挂载访问OSS的客户端,顺序读写能力强,充分发挥OSS高带宽优势。

ossfs 2.0适用于对存储访问性能要求较高的场景,比如AI训练、推理、大数据处理、自动驾驶等新型计算密集型负载。这类工作负载主要涉及顺序和随机读取、顺序(仅支持追加)写入操作,并且无需依赖完整的 POSIX 语义。

OSS Connector for AI/ML

PAI集成了OSS Connector for AI/ML,可在PyTorch代码中直接流式读取OSS文件。优势如下:

  • 流式加载:无需提前下载数据,节省GPU等待时间和成本。

  • 接口友好:与PyTorch Dataset用法一致,简单易用,封装优于OSS SDK,便于自定义。

  • 高效读取:相比OSS SDK,数据读取性能更优,加载更快。

该方式无需挂载OSS。如果您用PyTorch训练,需要读取海量(百万级别)小文件且对吞吐量要求较高,可使用OSS Connector for AI/ML加速读取。

OSS SDK

使用OSS2流式访问OSS数据。OSS2灵活高效,可缩短请求时间,提升训练效率。

如果您只需临时、非挂载地访问OSS,或按业务逻辑决定是否访问OSS,可使用OSS Python SDKAPI。

重要

通过JindoFuseossfs 2.0挂载OSS时,可指定访问OSSRAM角色,获得更精细的权限管控。同时支持在工作空间通用配置中开启OSS挂载必选RAM角色。开启后:

  • DSW/DLC/EAS等功能新建OSS挂载时必须选择RAM角色,不允许使用默认RAM角色;

  • 新建高级型/逻辑型数据集且涉及OSS挂载时,也会强制校验RAM角色;

  • DSW创建实例、DLC创建任务、EAS创建服务时,RAM角色不允许选择PAI默认角色。

JindoFuse

DLCDSW支持使用JindoFuseOSS数据集或路径挂载到容器,训练时直接读写OSS数据。

挂载方式

DLC中挂载OSS

创建分布式训练(DLC)任务时,可挂载OSS数据。挂载后,训练代码可以像访问本地文件一样读取OSS数据。支持以下两种挂载类型,具体配置请参见创建训练任务

数据集直接挂载区域,单击展开高级配置查看更多挂载选项,通过是否只读开关控制挂载模式。

挂载类型

描述

数据集

选择对象存储OSS类型的数据集,并配置挂载路径。公共数据集仅支持只读挂载。

直接挂载

直接挂载OSS Bucket存储路径。

使用已开启本地缓存的灵骏智算资源配额时,可开启使用缓存开关。

DSW中挂载OSS

创建DSW实例时,可挂载OSS数据。挂载后,开发环境可以像访问本地文件一样读取OSS数据。支持以下两种挂载类型,具体配置请参见创建DSW实例

两种挂载类型均提供展开高级配置选项。存储路径挂载支持OSS、通用型NAS、极速型NAS、CPFS和智算CPFS。

挂载类型

描述

数据集挂载

选择对象存储OSS类型的数据集,并配置挂载路径。公共数据集仅支持只读挂载。

存储路径挂载

直接挂载OSS Bucket存储路径。

默认配置限制

高级配置参数为空时使用默认配置,限制如下:

  • 为加速读取,挂载OSS时会缓存元数据(目录与文件列表)。

    分布式任务中,多个节点同时创建同一目录时,元数据缓存会让每个节点都尝试创建。最终只有一个节点成功,其余报错。

  • 默认使用OSSMultipart API创建文件,写入过程中OSS上看不到该对象,写入完成后才能查看。

  • 不支持同时读写同一文件。

  • 不支持随机写入文件。

常见JindoFuse配置

您也可以按场景在高级配置中自定义JindoFuse参数。

以下仅提供部分场景建议,未覆盖所有最优性能设置。更多配置请参见JindoFuse使用指南
  • 快速读写:允许用户读写,读取速度快,但并发读写可能会出现数据不一致的问题,适合挂载训练数据和模型,不适合作为工作目录。

    {
      "fs.oss.download.thread.concurrency": "cpu核数2倍",
      "fs.oss.upload.thread.concurrency": "cpu核数2倍",
      "fs.jindo.args": "-oattr_timeout=3 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink"
    }
    
  • 增量读写:在增量写入时能够保证数据一致性,覆盖原有数据会有一致性问题。读取速度略慢,适合保存训练的模型权重文件。

    {
      "fs.oss.upload.thread.concurrency": "cpu核数2倍",
      "fs.jindo.args": "-oattr_timeout=3 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink"
    }
    
  • 读写一致:在并发读写中能保持数据一致性,适用于对数据一致性要求高,可以容忍读取速度慢的场景,适合保存代码项目。

    {
      "fs.jindo.args": "-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink"
    }
    
  • 只读:仅允许读取,不允许写入,适合挂载公共数据集。

    {
      "fs.oss.download.thread.concurrency": "cpu核数2倍",
      "fs.jindo.args": "-oro -oattr_timeout=7200 -oentry_timeout=7200 -onegative_timeout=7200 -okernel_cache -ono_symlink"
    }

常见配置如下:

  • 选择JindoFuse版本

    {
      "fs.jindo.fuse.pod.image.tag": "6.7.0"
    }
  • 关闭元数据缓存:分布式任务中多个节点同时向同一目录写文件时,缓存可能导致部分节点写入失败。在JindoFuse命令行参数中增加-oattr_timeout=0-oentry_timeout=0-onegative_timeout=0可避免该问题。

    {
      "fs.jindo.args": "-oattr_timeout=0-oentry_timeout=0-onegative_timeout=0"
    }
  • 调整上传/下载线程数:配置以下参数调整并发。

    {
      "fs.oss.upload.thread.concurrency": "32",
      "fs.oss.download.thread.concurrency": "32",
      "fs.oss.read.readahead.buffer.count": "64",
      "fs.oss.read.readahead.buffer.size": "4194304"
    }
  • 使用AppendObject方式挂载OSS文件:本地创建的文件会调用OSS AppendObject接口。生成的Object大小不得超过5 GB,更多限制请参见AppendObject。示例配置如下。

    {
      "fs.jindo.args": "-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0",
      "fs.oss.append.enable": "true",
      "fs.oss.flush.interval.millisecond": "1000",
      "fs.oss.read.readahead.buffer.size": "4194304",
      "fs.oss.write.buffer.size": "262144"
    }
  • 挂载OSS-HDFS:开通OSS-HDFS请参见什么是OSS-HDFS服务。分布式训练场景建议增加以下参数。

    {
      "fs.jindo.args": "-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0 -ono_symlink -ono_xattr -ono_flock -odirect_io",
      "fs.oss.flush.interval.millisecond": "10000",
      "fs.oss.randomwrite.sync.interval.millisecond": "10000"
    }
  • 配置内存资源:通过fs.jindo.fuse.pod.mem.limit参数调整内存上限,示例如下。

    {
      "fs.jindo.fuse.pod.mem.limit": "10Gi"
    }

使用Python SDK修改数据集的JindoFuse参数

您还可以通过Python SDK修改JindoFuse参数。

  1. 完成以下准备工作。

    1. 安装工作空间SDK。

      !pip install alibabacloud-aiworkspace20210204
    2. 配置环境变量。环境变量用于SDK身份认证,避免在代码中硬编码密钥。具体步骤,请参见安装Credentials工具Linux、macOSWindows系统配置环境变量

  2. 修改JindoFuse参数。

    快速读写

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def change_config():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        # 建议值设置为CPU核数的2倍
        options['fs.oss.download.thread.concurrency'] = 32
        options['fs.oss.upload.thread.concurrency'] = 32
        options['fs.jindo.args'] = '-oattr_timeout=3 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink'
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    change_config()

    增量读写

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def change_config():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        # 建议值设置为CPU核数的2倍
        options['fs.oss.upload.thread.concurrency'] = 32
        options['fs.jindo.args'] = '-oattr_timeout=3 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink'
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    change_config()

    读写一致

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def change_config():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        options['fs.jindo.args'] = '-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0 -oauto_cache -ono_symlink'
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    change_config()

    只读

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def change_config():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        # 建议值设置为CPU核数的2倍
        options['fs.oss.download.thread.concurrency'] = 32
        options['fs.jindo.args'] = '-oro -oattr_timeout=7200 -oentry_timeout=7200 -onegative_timeout=7200 -okernel_cache -ono_symlink'
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    change_config()

    选择JindoFuse版本

    示例代码如下:

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def change_version():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        # 配置JindoFuse版本,可选6.4.4、6.7.0、6.6.0。release note见:https://aliyun.github.io/alibabacloud-jindodata/releases/
        options['fs.jindo.fuse.pod.image.tag'] = "6.7.0"
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    change_version()

    关闭元数据缓存

    分布式任务中多个节点同时向同一目录写文件时,Cache可能导致部分节点写入失败。在fuse命令行参数中增加-oattr_timeout=0-oentry_timeout=0-onegative_timeout=0可解决该问题。示例代码如下。

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def turnOffMetaCache():
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
        workspace_client = AIWorkspaceClient(
          config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
          )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
    
        options['fs.jindo.args'] = '-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0'
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    turnOffMetaCache()
    

    调整上传/下载线程数

    配置以下参数调整线程数:

    • fs.oss.upload.thread.concurrency:32

    • fs.oss.download.thread.concurrency:32

    • fs.oss.read.readahead.buffer.count:64

    • fs.oss.read.readahead.buffer.size:4194304

    示例代码如下:

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def adjustThreadNum():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
    
        options['fs.oss.upload.thread.concurrency'] = 32
        options['fs.oss.download.thread.concurrency'] = 32
        options['fs.oss.read.readahead.buffer.count'] = 32
     
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
     
     
    adjustThreadNum()
    

    使用AppendObject方式挂载OSS文件

    本地创建的文件会调用OSS AppendObject接口。生成的Object大小不得超过5 GB,更多限制请参见AppendObject。示例代码如下:

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def useAppendObject():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
    
        options['fs.jindo.args'] = '-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0'
        options['fs.oss.append.enable'] = "true"
        options['fs.oss.flush.interval.millisecond'] = "1000"
        options['fs.oss.read.buffer.size'] = "262144"
        options['fs.oss.write.buffer.size'] = "262144"
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    useAppendObject()

    挂载OSS-HDFS

    开通OSS-HDFS请参见什么是OSS-HDFS服务。使用OSS-HDFS Endpoint创建数据集的示例代码如下:

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import CreateDatasetRequest
    
    
    def createOssHdfsDataset():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        workspace_id = '** DLC任务所在工作空间ID **'
    
        oss_bucket = '** OSS-Bucket **'
        # OSS-HDFSEndpoint。
        oss_endpoint = f'{region_id}.oss-dls.aliyuncs.com'
        # 要挂载的OSS-HDFS路径。
        oss_path = '/'
        # 本地挂载路径。
        mount_path = '/mnt/data/'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
    
        response = workspace_client.create_dataset(CreateDatasetRequest(
            workspace_id=workspace_id,
            name="** 数据集的名字 **",
            data_type='COMMON',
            data_source_type='OSS',
            property='DIRECTORY',
            uri=f'oss://{oss_bucket}.{oss_endpoint}{oss_path}',
            accessibility='PRIVATE',
            source_type='USER',
            options=json.dumps({
                'mountPath': mount_path,
                # 分布式训练场景建议增加以下参数。
                'fs.jindo.args': '-oattr_timeout=0 -oentry_timeout=0 -onegative_timeout=0 -ono_symlink -ono_xattr -ono_flock -odirect_io',
                'fs.oss.flush.interval.millisecond': "10000",
                'fs.oss.randomwrite.sync.interval.millisecond': "10000",
            })
        ))
        print(f'datasetId: {response.body.dataset_id}')
    
    createOssHdfsDataset()
    
    

    配置内存资源

    通过fs.jindo.fuse.pod.mem.limit参数调整内存资源,示例代码如下:

    import json
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    from alibabacloud_aiworkspace20210204.client import Client as AIWorkspaceClient
    from alibabacloud_aiworkspace20210204.models import UpdateDatasetRequest
    
    
    def adjustResource():
        # 使用DLC任务所在地域。例如华东1(杭州)配置为cn-hangzhou。
        region_id = 'cn-hangzhou'
        # AccessKey拥有所有API访问权限,建议使用RAM用户。
        # 请勿将AccessKey IDAccessKey Secret保存到工程代码中,避免泄露。
        # 本示例通过Credentials SDK从环境变量读取AccessKey。请提前安装Credentials工具并配置环境变量。
        cred = CredClient()
        dataset_id = '** 数据集的ID **'
    
        workspace_client = AIWorkspaceClient(
            config=Config(
                credential=cred,
                region_id=region_id,
                endpoint="aiworkspace.{}.aliyuncs.com".format(region_id),
            )
        )
        # 1、获取数据集内容
        get_dataset_resp = workspace_client.get_dataset(dataset_id)
        options = json.loads(get_dataset_resp.body.options)
        # 要配置的内存资源。
        options['fs.jindo.fuse.pod.mem.limit'] = "10Gi"
    
        update_request = UpdateDatasetRequest(
            options=json.dumps(options)
        )
        # 2、更新options
        workspace_client.update_dataset(dataset_id, update_request)
        print('new options is: {}'.format(update_request.options))
    
    
    adjustResource()
    

ossfs 2.0

挂载OSS数据源时,在高级配置中设置{"mountType":"ossfs"},即可使用ossfs挂载。

挂载方式

DLC中挂载OSS

创建分布式训练(DLC)任务时,可挂载OSS数据。挂载后,训练代码可以像访问本地文件一样读取OSS数据。支持以下两种挂载类型,具体配置请参见创建训练任务

挂载类型

描述

数据集

选择对象存储OSS类型的数据集,并配置挂载路径。公共数据集仅支持只读挂载。

直接挂载

直接挂载OSS Bucket存储路径。

使用已开启本地缓存的灵骏智算资源配额时,可开启使用缓存开关。

DSW中挂载OSS

创建DSW实例时,可挂载OSS数据。挂载后,开发环境可以像访问本地文件一样读取OSS数据。支持以下两种挂载类型,具体配置请参见创建DSW实例

挂载类型

描述

数据集挂载

选择对象存储OSS类型的数据集,并配置挂载路径。公共数据集仅支持只读挂载。

存储路径挂载

直接挂载OSS Bucket存储路径。

常见ossfs配置

高级配置中,通过fs.ossfs.args设置高级参数,多个参数之间使用半角逗号,分隔。高级参数说明请参见ossfs 2.0。常见场景示例如下:

  • 任务过程中数据源不变:如果读取过程中文件不会被修改,可配置较大的缓存时间,减少元数据请求次数。例如读取一批已有旧文件,处理后生成新文件。

    {
        "mountType":"ossfs",
        "fs.ossfs.args": "-oattr_timeout=7200" 
    }
  • 快速读写:使用较小的元数据缓存时间,平衡缓存效率与数据及时性。

    {
        "mountType":"ossfs",
        "fs.ossfs.args": "-oattr_timeout=3, -onegative_timeout=0"
    }
  • 分布式任务读写一致:ossfs默认会缓存元数据。通过以下配置实现多节点同步视图。

    {   
        "mountType":"ossfs",
        "fs.ossfs.args": "-onegative_timeout=0, -oclose_to_open"
    }
  • DLC/DSW场景中同时打开过多文件导致OOM:DLCDSW场景中任务并发量较高,可能同时打开大量文件导致OOM(内存不足)。通过以下配置缓解内存压力。

    {
        "mountType":"ossfs",
        "fs.ossfs.args": "-oreaddirplus=false, -oinode_cache_eviction_threshold=300000"
    }
  • 写入大文件失败-oupload_buffer_size用于设置分片上传缓冲区大小(Bytes)。该参数决定可写入文件的最大大小,计算方式为upload_buffer_size * 10000。

    ossfs 2.0默认分片大小为8 MiB,最大支持写入78.125 GiB的文件。超过该限制时写入会失败。可配置-oupload_buffer_size增加分片大小,从而提升可写入文件大小上限。例如将Part大小设为32 MiB(33554432字节),可支持最大312.5 GiB的文件。-upload_buffer_size越大,消耗内存越多,可配置-total_mem_limit控制内存使用,详情请参见挂载选项说明

    {
        "mountType":"ossfs",
        "fs.ossfs.args": "-oupload_buffer_size=33554432"
    }
    

OSS Connector for AI/ML

使用OSS Connector for AI/ML加速模型训练是阿里云OSS团队为AI/ML场景设计的客户端库,在大规模PyTorch训练中提供便捷的数据加载,缩短数据传输时间,加速模型训练。PAI集成了OSS Connector for AI/ML,可在PyTorch代码中流式读取OSS文件。

使用限制

  • 官方镜像:仅当DLC任务或DSW实例选择PyTorch 2.0及以上版本的镜像时,可用OSS Connector for AI/ML模块。

  • 自定义镜像:仅支持PyTorch 2.0及以上版本。满足版本要求后,运行以下命令安装OSS Connector for AI/ML模块。

    pip install -i http://yum.tbsite.net/aliyun-pypi/simple/ --extra-index-url http://yum.tbsite.net/pypi/simple/ --trusted-host=yum.tbsite.net osstorchconnector
  • Python版本:仅支持Python 3.8~3.12。

准备工作

  1. 配置credential文件。

    通过以下任一方式配置credential:

    • 参考配置DLC RAM角色,为DLC任务配置免密访问OSScredential。DLC任务会获取STS临时访问凭证,安全访问OSS或其他云资源,无需在代码中配置认证信息,降低密钥泄露风险。

    • 在代码项目中配置credential文件管理认证信息。配置示例如下:

      说明

      明文配置AK存在安全风险,建议在DLC实例内通过角色自动配置credential,详情请参见配置DLC RAM角色

      使用OSS Connector for AI/ML接口时,指定credential文件路径后自动获取认证信息,完成OSS数据请求认证。

      {
        "AccessKeyId": "<Access-key-id>",
        "AccessKeySecret": "<Access-key-secret>",
        "SecurityToken": "<Security-Token>",
        "Expiration": "2024-08-20T00:00:00Z"
      }

      具体配置项说明如下:

      配置项

      是否必填

      说明

      示例值

      AccessKeyId

      阿里云账号或RAM用户的AccessKey IDAccessKey Secret。

      说明

      使用STS临时访问凭证访问OSS时,请设置为临时访问凭证的AccessKey IDAccessKey Secret。

      NTS****

      AccessKeySecret

      7NR2****

      SecurityToken

      临时访问令牌。使用STS临时访问凭证访问OSS时,需要设置。

      STS.6MC2****

      Expiration

      鉴权信息过期时间,Expiration为空表示永不过期,过期后OSS Connector会重新读取鉴权信息。

      2024-08-20T00:00:00Z

  2. 配置config.json文件,示例如下:

    在代码项目中配置config.json文件,管理并发数、预取参数和日志路径等。使用OSS Connector for AI/ML接口时,指定config.json文件路径后,系统会自动读取并发数和预取值,并将OSS数据请求日志输出到指定文件。

    {
        "logLevel": 1,
        "logPath": "/var/log/oss-connector/connector.log",
        "auditPath": "/var/log/oss-connector/audit.log",
        "datasetConfig": {
            "prefetchConcurrency": 24,
            "prefetchWorker": 2
        },
        "checkpointConfig": {
            "prefetchConcurrency": 24,
            "prefetchWorker": 4,
            "uploadConcurrency": 64
        }
    }

    配置项说明如下:

    配置项

    是否必填

    说明

    示例值

    logLevel

    日志记录级别。默认为INFO。取值如下:

    • 0:表示Debug。

    • 1:表示INFO。

    • 2:表示WARN。

    • 3:表示ERROR。

    1

    logPath

    connector日志路径。默认路径为/var/log/oss-connector/connector.log

    /var/log/oss-connector/connector.log

    auditPath

    connector IO审计日志,记录延迟大于100毫秒的读写请求。默认路径为/var/log/oss-connector/audit.log

    /var/log/oss-connector/audit.log

    DatasetConfig

    prefetchConcurrency

    使用DatasetOSS预取数据时的并发数,默认为24。

    24

    prefetchWorker

    使用DatasetOSS预取时可用的vCPU数,默认为4。

    2

    checkpointConfig

    prefetchConcurrency

    使用checkpoint readOSS预取数据时的并发数,默认为24。

    24

    prefetchWorker

    使用checkpoint readOSS预取时可用的vCPU数,默认为4。

    4

    uploadConcurrency

    使用checkpoint write上传数据时的并发数,默认为64。

    64

使用方式

OSS Connector for AI/ML提供OssMapDatasetOssIterableDataset两种数据集访问接口,分别扩展DatasetIterableDataset。OssIterableDataset有预取优化,训练效率更高。OssMapDataset的读取顺序由DataLoader决定,支持shuffle。建议如下:

  • 如果内存较小或数据量较大,只需顺序读取且对并行处理要求不高,建议使用OssIterableDataset。

  • 如果内存充足、数据量较小,且需要随机操作和并行处理,建议使用OssMapDataset。

OSS Connector for AI/ML还提供OssCheckpoint接口,支持模型加载和保存。当前OssCheckpoint功能仅限通用资源环境使用。

下面介绍这三种接口的使用方式:

OssMapDataset

支持以下三种数据集访问模式:

  • 根据OSS路径前缀访问文件夹

    只需指定文件夹名称,无需配置索引文件。如果OSS文件夹结构如下,可采用该方式访问数据集:

    dataset_folder/
        ├── class1/
        │   ├── image1.JPEG
        │   └── ...
        ├── class2/
        │   ├── image2.JPEG
        │   └── ...

    使用时指定OSS路径前缀,并自定义文件流解析方式。以下是解析和转换图片文件的方法:

    def read_and_transform(data):
        normalize = transforms.Normalize(mean=[0.485, 0.456, 0.406],
                                         std=[0.229, 0.224, 0.225])
        transform = transforms.Compose([
            transforms.RandomResizedCrop(224),
            transforms.RandomHorizontalFlip(),
            transforms.ToTensor(),
            normalize,
        ])
    
        try:
            img = accimage.Image((data.read()))
            val = transform(img)
            label = data.label # 文件名
        except Exception as e:
            print("read failed", e)
            return None, 0
        return val, label
    dataset = OssMapDataset.from_prefix("{oss_data_folder_uri}", endpoint="{oss_endpoint}", transform=read_and_transform, cred_path=cred_path, config_path=config_path)
  • 根据manifest_file获取文件

    支持访问多个OSS Bucket,管理更灵活。如果OSS文件夹结构如下,且存在管理文件名与Label对应关系的manifest_file,可采用该方式访问数据集。

    dataset_folder/
        ├── class1/
        │   ├── image1.JPEG
        │   └── ...
        ├── class2/
        │   ├── image2.JPEG
        │   └── ...
        └── .manifest

    其中manifest_file格式如下:

    {'data': {'source': 'oss://examplebucket.oss-cn-wulanchabu.aliyuncs.com/dataset_folder/class1/image1.JPEG'}}
    {'data': {'source': ''}}

    使用时自定义manifest_file的解析方式,示例如下:

    def transform_oss_path(input_path):
        pattern = r'oss://(.*?)\.(.*?)/(.*)'
        match = re.match(pattern, input_path)
        if match:
            return f'oss://{match.group(1)}/{match.group(3)}'
        else:
            return input_path
    
    
    def manifest_parser(reader: io.IOBase) -> Iterable[Tuple[str, str, int]]:
        lines = reader.read().decode("utf-8").strip().split("\n")
        data_list = []
        for i, line in enumerate(lines):
            data = json.loads(line)
            yield transform_oss_path(data["data"]["source"]), ""
    dataset = OssMapDataset.from_manifest_file("{manifest_file_path}", manifest_parser, "", endpoint=endpoint, transform=read_and_trans, cred_path=cred_path, config_path=config_path)
  • 根据OSS_URI列表的方式获取文件

    只需指定OSS URI,无需配置索引文件即可访问OSS文件。示例如下:

    uris =["oss://examplebucket.oss-cn-wulanchabu.aliyuncs.com/dataset_folder/class1/image1.JPEG", "oss://examplebucket.oss-cn-wulanchabu.aliyuncs.com/dataset_folder/class2/image2.JPEG"]
    dataset = OssMapDataset.from_objects(uris, endpoint=endpoint, transform=read_and_trans, cred_path=cred_path, config_path=config_path)

OssIterableDataset

OssIterableDataset也支持三种访问方式,与OssMapDataset相同。用法如下:

  • 根据OSS路径前缀访问文件夹

    dataset = OssIterableDataset.from_prefix("{oss_data_folder_uri}", endpoint="{oss_endpoint}", transform=read_and_transform, cred_path=cred_path, config_path=config_path)
  • 根据manifest_file获取文件

    dataset = OssIterableDataset.from_manifest_file("{manifest_file_path}", manifest_parser, "", endpoint=endpoint, transform=read_and_trans, cred_path=cred_path, config_path=config_path)
  • 根据OSS_URI列表的方式获取文件

    dataset = OssIterableDataset.from_objects(uris, endpoint=endpoint, transform=read_and_trans, cred_path=cred_path, config_path=config_path)

OssCheckpoint

OssCheckpoint功能仅支持通用计算资源环境。OSS Connector for AI/ML支持通过OssCheckpoint读取和保存OSS模型文件。用法如下:

checkpoint = OssCheckpoint(endpoint="{oss_endpoint}", cred_path=cred_path, config_path=config_path)

checkpoint_read_uri = "{checkpoint_path}"
checkpoint_write_uri = "{checkpoint_path}"
with checkpoint.reader(checkpoint_read_uri) as reader:
    state_dict = torch.load(reader)
    model.load_state_dict(state_dict)
with checkpoint.writer(checkpoint_write_uri) as writer:
    torch.save(model.state_dict(), writer)

代码示例

以下是OSS Connector for AI/ML的示例代码,可直接用于访问OSS数据:

from osstorchconnector import OssMapDataset, OssCheckpoint
import torchvision.transforms as transforms
import accimage
import torchvision.models as models
import torch

cred_path = "/mnt/.alibabacloud/credentials"  # 为DLC任务和DSW实例配置角色后credential的默认路径。
config_path = "config.json"
checkpoint = OssCheckpoint(endpoint="{oss_endpoint}", cred_path=cred_path, config_path=config_path)
model = models.__dict__["resnet18"]()

epochs = 100  # 指定epoch
checkpoint_read_uri = "{checkpoint_path}"
checkpoint_write_uri = "{checkpoint_path}"
with checkpoint.reader(checkpoint_read_uri) as reader:
    state_dict = torch.load(reader)
    model.load_state_dict(state_dict)


def read_and_transform(data):
    normalize = transforms.Normalize(mean=[0.485, 0.456, 0.406],
                                     std=[0.229, 0.224, 0.225])
    transform = transforms.Compose([
        transforms.RandomResizedCrop(224),
        transforms.RandomHorizontalFlip(),
        transforms.ToTensor(),
        normalize,
    ])

    try:
        img = accimage.Image((data.read()))
        value = transform(img)
    except Exception as e:
        print("read failed", e)
        return None, 0
    return value, 0
dataset = OssMapDataset.from_prefix("{oss_data_folder_uri}", endpoint="{oss_endpoint}", transform=read_and_transform, cred_path=cred_path, config_path=config_path)
data_loader = torch.utils.data.DataLoader(
    dataset, batch_size="{batch_size}",num_workers="{num_workers"}, pin_memory=True)

for epoch in range(args.epochs):
    for step, (images, target) in enumerate(data_loader):
        # batch processing
        # model training
    # save model
    with checkpoint.writer(checkpoint_write_uri) as writer:
        torch.save(model.state_dict(), writer)

上述代码的关键实现说明如下:

  • OssMapDataset直接基于给定的OSS URI,构建与PyTorch DataLoader用法一致的dataset。

  • 用该dataset构建标准DataLoader,并循环DataLoader执行训练流程,如处理当前batch、模型训练与保存等。

  • 该过程无需将数据集挂载到容器,也无需预先将数据存到本地,数据按需加载。

OSS SDK

OSS Python SDK

您可直接使用OSS Python SDK读写OSS数据,操作步骤如下:

  1. 安装OSS Python SDK。详情请参见安装(Python SDK V1)

  2. OSS Python SDK配置访问凭证。凭证用于验证您对OSS的访问权限,详情请参见配置访问凭证(Python SDK V1)

  3. 读写OSS数据。

    # -*- coding: utf-8 -*-
    import oss2
    from oss2.credentials import EnvironmentVariableCredentialsProvider
    
    # 使用环境变量中的RAM用户访问密钥配置访问凭证
    auth = oss2.ProviderAuth(EnvironmentVariableCredentialsProvider())
    bucket = oss2.Bucket(auth, '<Endpoint>', '<your_bucket_name>')
    # 读取完整文件。
    result = bucket.get_object('<your_file_path/your_file>')
    print(result.read())
    # 按Range读取数据。
    result = bucket.get_object('<your_file_path/your_file>', byte_range=(0, 99))
    # 写数据到OSS。
    bucket.put_object('<your_file_path/your_file>', '<your_object_content>')
    # 对Appendable类型文件追加数据。
    result = bucket.append_object('<your_file_path/your_file>', 0, '<your_object_content>')
    result = bucket.append_object('<your_file_path/your_file>', result.next_position, '<your_object_content>')
    

    请按实际修改以下配置项:

    配置项

    描述

    <Endpoint>

    填写Bucket所在地域的Endpoint。以华东1(杭州)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。获取Endpoint更多信息,请参见地域和Endpoint

    <your_bucket_name>

    填写存储空间名称。

    <your_file_path/your_file>

    待读写的文件路径。填写Object完整路径,不包含Bucket名称,例如testfolder/exampleobject.txt

    <your_object_content>

    Append的内容,按实际情况修改。

OSS Python API

使用OSS Python API可方便存储训练数据和模型。开始前确保已安装OSS Python SDK并正确设置访问凭证,详情请参见安装(Python SDK V1)配置访问凭证(Python SDK V1)

  • 加载训练数据

    您可将数据存放在一个OSS Bucket中,将数据路径与对应Label存入同一OSS Bucket的索引文件。通过自定义Dataset,在PyTorch中用DataLoaderAPI多进程并行读取数据,可提升训练效率,示例如下。

    import io
    import oss2
    from oss2.credentials import EnvironmentVariableCredentialsProvider
    import PIL
    import torch
    
    class OSSDataset(torch.utils.data.dataset.Dataset):
        def __init__(self, endpoint, bucket, auth, index_file):
            self._bucket = oss2.Bucket(auth, endpoint, bucket)
            self._indices = self._bucket.get_object(index_file).read().split(',')
    
        def __len__(self):
            return len(self._indices)
    
        def __getitem__(self, index):
            img_path, label = self._indices(index).strip().split(':')
            img_str = self._bucket.get_object(img_path)
            img_buf = io.BytesIO()
            img_buf.write(img_str.read())
            img_buf.seek(0)
            img = Image.open(img_buf).convert('RGB')
            img_buf.close()
            return img, label
    
    
    # 从环境变量获取访问凭证。运行本代码前,请确保已设置环境变量OSS_ACCESS_KEY_IDOSS_ACCESS_KEY_SECRET。
    auth = oss2.ProviderAuth(EnvironmentVariableCredentialsProvider())
    dataset = OSSDataset(endpoint, bucket, auth, index_file)
    data_loader = torch.utils.data.DataLoader(
        dataset,
        batch_size=batch_size,
        num_workers=num_loaders,
        pin_memory=True)
    

    关键配置说明如下:

    关键配置

    描述

    endpoint

    填写Bucket所在地域的Endpoint。以华东1(杭州)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。获取Endpoint更多信息,请参见地域和Endpoint

    bucket

    填写Bucket名称。

    index_file

    索引文件的路径。

    说明

    示例中索引文件格式为每条样本用英文逗号(,)分隔,样本路径与Label之间用英文冒号(:)分隔。

  • SaveLoad模型

    您可使用OSS Python API保存或加载PyTorch模型(PyTorch保存/加载模型详情请参见PyTorch),示例如下:

    • Save模型

      from io import BytesIO
      import torch
      import oss2
      from oss2.credentials import EnvironmentVariableCredentialsProvider
      
      auth = oss2.ProviderAuth(EnvironmentVariableCredentialsProvider())
      # Bucket名称
      bucket_name = "<your_bucket_name>"
      bucket = oss2.Bucket(auth, endpoint, bucket_name)
      buffer = BytesIO()
      torch.save(model.state_dict(), buffer)
      bucket.put_object("<your_model_path>", buffer.getvalue())
      

      其中

      • endpointBucket所在地域的Endpoint。以华东1(杭州)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。

      • <your_bucket_name>OSS Bucket名称,且开头不带oss://

      • <your_model_path>为模型路径,按实际情况修改。

    • Load模型

      from io import BytesIO
      import torch
      import oss2
      from oss2.credentials import EnvironmentVariableCredentialsProvider
      
      auth = oss2.ProviderAuth(EnvironmentVariableCredentialsProvider())
      bucket_name = "<your_bucket_name>"
      bucket = oss2.Bucket(auth, endpoint, bucket_name)
      buffer = BytesIO(bucket.get_object("<your_model_path>").read())
      model.load_state_dict(torch.load(buffer))

      其中

      • endpointBucket所在地域的Endpoint。以华东1(杭州)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。

      • <your_bucket_name>OSS Bucket名称,且开头不带oss://

      • <your_model_path>为模型路径,按实际情况修改。