Fuxi: a Fault-Tolerant Resource Management and Job Scheduling System at Internet Scale

Fuxi: a Fault-Tolerant Resource Management and Job Scheduling System at Internet Scale
复制标题

DOI:
10.14778/2733004.2733012
复制
发表时间:
2014-08
期刊:
Proc. VLDB Endow.
影响因子:
--
通讯作者:
Zhuo Zhang;C. Li;Y. Tao;Renyu Yang;Hong Tang;Jie Xu
Zhuo Zhang;C. Li;Y. Tao;Renyu Yang;Hong Tang;Jie Xu
中科院分区:
其他
文献类型:
--
作者:
Zhuo Zhang;C. Li;Y. Tao;Renyu Yang;Hong Tang;Jie Xu

文献摘要

被引文献

相似文献

可伸缩性和容错是互联网规模的分布式计算面临的两个基本挑战。尽管学术界和工业界最近都取得了许多进展,但这两个问题仍远未解决。在本文中,我们介绍了Fuxi,一个资源管理和工作调度系统,能够处理阿里巴巴的工作量,每天生成和分析数百tb的数据,以帮助优化公司的业务运营和用户体验。我们采用了几项新技术,使伏羲能够在拥有数千个节点的大型集群上执行数十万并发任务的高效调度:1)支持多维资源分配和数据局部性的增量资源管理协议;2)用户透明的故障恢复,任何福喜组件的故障都不会影响用户作业的执行;3)有效的检测机制和多级黑名单方案,防止其影响作业执行。我们的评估结果表明,在合成工作负载下,可以实现95%和91%的调度CPU/内存利用率,并且Fuxi能够在GraySort中实现2.36T-B/分钟的吞吐量。此外,在5%的故障注入率下,同样的伏西作业只经历了大约16%的减速。当我们将故障注入率提高一倍至10%时,减速速度仅增长到20%。自2009年以来,Fuxi已经部署在我们的生产环境中,现在它管理着数十万个服务器节点。
Scalability and fault-tolerance are two fundamental challenges for all distributed computing at Internet scale. Despite many recent advances from both academia and industry, these two problems are still far from settled. In this paper, we present Fuxi, a resource management and job scheduling system that is capable of handling the kind of workload at Alibaba where hundreds of terabytes of data are generated and analyzed everyday to help optimize the company's business operations and user experiences. We employ several novel techniques to enable Fuxi to perform efficient scheduling of hundreds of thousands of concurrent tasks over large clusters with thousands of nodes: 1) an incremental resource management protocol that supports multi-dimensional resource allocation and data locality; 2) user-transparent failure recovery where failures of any Fuxi components will not impact the execution of user jobs; and 3) an effective detection mechanism and a multi-level blacklisting scheme that prevents them from affecting job execution. Our evaluation results demonstrate that 95% and 91% scheduled CPU/memory utilization can be fulfilled under synthetic workloads, and Fuxi is capable of achieving 2.36T-B/minute throughput in GraySort. Additionally, the same Fuxi job only experiences approximately 16% slowdown under a 5% fault-injection rate. The slowdown only grows to 20% when we double the fault-injection rate to 10%. Fuxi has been deployed in our production environment since 2009, and it now manages hundreds of thousands of server nodes.