<span class=“js_title_inner“>技术博客|第19期:Hulu分布式训练平台的架构与实践(上)</span>

目前机器学习已应用在Hulu的多个场景,包括推荐、搜索、广告、异常登录检测、语音识别等领域。作为机器学习生命周期的重要一环,Hulu机器学习平台为数据科学家和算法工程师提供了涵盖算力、通信、调度、跟踪、调优的工具和能力,以便更好、更快、更有效的让模型带来商业价值。本文将着重介绍Hulu分布式训练平台的架构,以及在平台搭建和维护中遇到的挑战和解决方案。

Hulu机器学习训练平台已经有五年的历史。五年中,平台围绕机器学习场景,整合了TensorFlow,LightGBM,PyTorch,XGBoost等等一系列机器学习框架,为搜索,推荐,广告,内容发现等等各个领域提供了算力和框架支撑。在平台的不断迭代中,随着公司业务规模的扩大,订阅用户数的增多,机器学习模型的复杂化,训练平台也在面临着方方面面的技术挑战。

随着样本数据、特征和标签数据集愈加庞大,单机训练已无法满足训练时长的要求,这将降低模型的更新频率,有可能导致模型推理效果变差。我们可以通过使用更高性能的GPU来缓解这个问题,但是仍然突破不了单卡训练的瓶颈。

模型参数一般为稀疏或稠密的向量或矩阵,通常模型中更多的参数可以表征更多的特性。当前深度学习模型中参数规模在以倍数速度增长,部分推荐模型已达到10TB级别,在这种情况下,单机/单卡内存已经远远不能满足需求。
为了解决上述问题,分布式训练平台应运而生。通过对分布式训练的支持,平台不仅实现了对关键模型的加速,而且也兼容了目前超大模型的通用解决方案。


目前,业内常见的分布式训练方案分为数据并行、张量并行和流水线并行等,从部署模式上又分为单机多卡模式和多机多卡模式等。单机多卡指作业的多个工作进程,在同一台物理机的多个GPU上训练,通过共享内存、Unix Socket或NVLink进行多进程之间的同步和通信,其特点是通信吞吐量大、系统简单,但受限于显存和算力,只适用于数据和模型有限规模下的场景。与单机多卡相比,多机模式突破了硬件对数据集和模型大小的限制,它将多个工作进程分配在多台物理机上,通过网络进行参数的同步。多机模型还可以利用集群的碎片和深度个性化配置来提高资源使用率。因此,平台需要同时支持这单机和多机这两种部署模式。
对于数据并行的分布式训练来说,参数同步的一个思路是通过参数服务器实现。参数服务器是机器学习训练的范式之一,其优势在于通过异步训练,解决机器学习过程中的掉队者问题,降低了GPU资源的一致性要求。我们还可以将分片逻辑融入参数服务器设计,从而打破单机内存对于海量模型参数的限制。但该模式会带来单点问题,在训练过程中,参数服务器的带宽成为系统的瓶颈。如果使用多个参数服务器,则通信模式变为“All-to-All”,这可能使网络饱和。因此该模式受限于机房的带宽,对部分大规模参数模型仍有一定局限性。
以Ring AllReduce为代表的参数同步算法,解决了训练过程中单点通信的问题,同参数服务器模式相比,它通过去中心节点化,降低了多机之间网络带宽压力,从而减少了梯度同步中的瓶颈,加快了训练速度,但对GPU资源的一致性有较高要求。本文将在关键问题和解决方案部分对Ring AllReduce算法展开介绍。综上,平台需要同时支持参数服务器和Ring Allreduce两种训练模式。
对于张量并行化、流水线并行化等模型并行方案,平台也需要提供类似数据并行方案中多机分布的基础设施层支持。和基于Ring Allreduce的分布式方案不同,模型对GPU资源的一致性要求并不严苛,因此提供灵活的硬件支持也是平台的一个重要功能。当用户使用第三方框架对模型或流程切分后,将利用平台个性化的多类型资源和高性能通信,以实现训练效率的提升。关于模型并行的应用层解决方案,后续将另择机重点讨论。
为了支持上述业务需求,平台需要解决资源抽象、任务调度、任务通信、应用接入等等一系列问题。

Hulu分布式训练平台通过整合硬件资源提供的算力、存储和通信设施,提供了端到端的、高性能的模型训练、模型存储与管理、模型加速、模型开发、模型测试和性能调优能力。通过对分布式训练的支持和硬件的虚拟抽象,平台解决了海量数据增长、模型结构复杂带来的训练速度低下的问题,实现了对于算力资源的高效利用,有效保障了Hulu多个场景下的业务需求。下面将对平台架构和各个逻辑层次做展开介绍。

Hulu的训练平台可以按照功能垂直分为3层,分别为应用接入层、基础设施层、硬件资源层。其中:
硬件资源层主要为平台提供资源保障与成本控制的能力。资源保障覆盖了算力、存储、通信等各个方面。成本控制则包含数据中心的容量规划和成本优化两部分。
基础设施层的使命是对资源进行抽象与包装,赋能上层的业务应用。在Kubernetes集群中,基础设施层提供了分布式调度、高性能容器网络、分布式共享存储、镜像、多租户管理以及日志监控等等诸多能力。
接入层主要提供业务落地时的各种支持能力。比如TensorFlow/Pytorch等框架接入,分布式训练支持,参数服务器支持,以及一系列开发、调试、性能剖析的能力。
接下来对平台架构各个层次展开介绍。

目前,Hulu机器学习训练平台的算力主要来源是自建机房的物理机。根据类型和用途的不同,可以将物理资源分为三类角色:GPU算力,CPU算力和存储单元。这其中:
GPU算力采用了Nvidia提供的DGX方案,绝大部分物理机拥有Ampere架构的GPU,少部分拥有Volta架构的GPU。物理机之间通过50Gbps光纤相连。GPU算力主要承载了较重的模型训练任务。
CPU算力采用了Intel Xeon(R) Gold方案,物理机之间通过20Gbps光纤相连。主要承载了数据准备,清洗和部分较轻的模型训练任务。
平台利用存储单元部署了Ceph集群,作为共享文件系统提供给分布式训练任务。
除了对平台的资源保障,硬件资源层还负责对数据中心的平台成本进行管理。
硬件资源层会定期根据平台上游团队的未来业务量进行估计,根据需求增长采购硬件;同时,硬件资源层还着眼业界硬件新架构、新特性来综合考虑硬件采购方案。
成本优化也是硬件资源层的重要使命之一。通过对资源的细粒度抽象和隔离,硬件资源层提高了平台的资源有效性。这在下文将继续讨论。

训练平台在基础设施层采用了基于Kubernetes的解决方案。基于Kubernetes强大的容器编排能力,丰富的自定义接口和活跃的云原生社区,基础设施层将底层资源的能力进行了抽象和包装, 为训练平台提供了强大的高性能算力、灵活的资源调度、高性能容器网络、可靠的配额管理、共享存储和监控/调试等等一系列能力。具体来说,基础设施层对训练任务提供了如下能力:
面向分布式、高性能计算的批调度器Volcano。一方面Volcano的批调度模式可以极大程度上避免资源闲置与浪费,另一方面Volcano可以与网络拓扑结合,给出能带来最大通信性能的负载放置。
高性能的容器网络方案Cilium。作为云原生社区火热的网络解决方案之一,Cilium不仅基于内核ebpf技术实现了极高的容器网络性能,而且对于不同内核版本,OS发行版都有广泛的兼容性。
分布式共享存储Ceph。基础设施层利用物理机空闲的存储能力构建了共享存储解决方案,以便于分布式作业的数据分片和加载。这套方案与计算集群混部,可以结合调度器实现计算靠近数据的调度策略,节省网络通信流量。
日志/监控平台。基础设施层为其上运行的训练任务提供了日志保存与检索,监控数据收集与展示功能。为研究员和算法工程师的任务排错,任务调试,任务调优提供了充分支持。
此外,基础设施层还依托于Kubernetes强大的配额管理能力和丰富的云原生组件,为用户提供了多租户和容器镜像管理的支持。

应用接入层作为将平台能力投递给用户的最后一公里,不仅承载着融合包装平台能力的功能,还肩负着简化用户开发,提升用户体验的责任。应用接入层不仅整合了TensorFlow、Pytorch、LightGBM等主流框架,而且为模型开发与调试、模型加速、模型性能剖析、模型部署全流程提供了完整的支持。具体来说,应用接入层包含以下能力:
分布式训练的支持。应用接入层通过集成Horovod等框架为研究员和算法工程师提供了高性能、易使用的分布式支持。
机器学习应用开发工具。应用接入层整合了大数据SDK,Ray,Jupyter Notebook和Tensorboard等开源工具,依托于基础设施层共享存储的能力,全方位支持用户的开发和调试过程。
模型部署工具。应用接入层通过集成Hulu模型市场的SDK,简化了模型打包、模型部署等流程,极大地节省了用户的时间成本。

以Horovod+TensorFlow1x为例,本文在此简述将单机任务以数据并行方式迁移到分布式平台所需要的改动:
单机模型改造
参考Horovod官方指导,在模型中使用hvd.init()初始化,并使用hvd.rank()作为序号进行数据加载;
使用hvd.DistributedOptimizer包装原有的TensorFlow优化器。
分布式任务提交
对于开发阶段的分布式任务,平台暴露了HTTP接口,并提供了一系列SDK方便用户提交任务到计算集群;
对于生产环境的分布式任务,平台提供了Airflow Operator和一系列的CICD工具,方便任务被Airflow托管,并周期性执行。
任务追踪。平台提供了一系列支持组件,来追踪任务的执行状态:
平台UI,用于展示任务状态和相关事件,并提供远程调试终端;
日志系统。平台基于Filebeat+ElasticSearch构建了日志系统,用于收集、存储任务日志,并提供日志搜索功能;
监控系统。平台基于Prometheus+Grafana构建了容器监控系统,用于收集和展示任务过程中对于各种硬件资源的使用情况。
除此之外,平台还在任务镜像中内置了诸如Jupyter等一系列工具,方便用户在容器内对任务进行调试。通过平台提供的强大基础设施支撑和丰富工具保障,用户的迁移工作得到了极大地简化。

本文阐述了Hulu分布式训练平台的项目背景、平台总体架构和各个逻辑层次。并且,本文给出了一个快速上手指南,方便读者对平台的使用方式有初步印象。在下篇中,本文将对实际平台搭建和维护中,出现的通信、调度、资源虚拟化等关键问题和解决方案进行深入剖析。

Chenyu Zheng,内容发现部门高级研发工程师。
Chengmin Yang,内容发现部门资深研发工程师。
内容发现部门(Content Discovery Org.)是迪士尼流媒体核心研发部门,主攻Hulu、Disney+、Star+等迪士尼流媒体产品线的三大业务方向:搜索、个性化推荐、内容推广。在每个业务方向上都和人工智能技术深度融合,涉及AI平台的搭建、前沿算法的研究、以及工程系统的集成,致力于为迪士尼流媒体用户提供最佳的视频观看体验。
自成立伊始,该部门始终将内容的精准传递作为首要业务目标,深入结合工程、算法和数据,利用人才优势与人工智能基础解决业务问题。

更多推荐



所有评论(0)