IP Library Granted Patent US 12,223,349
Granted Patent B2
US 12,223,349 · App. 17/379,742 · Granted Feb 11, 2025

Utilization-aware resource scheduling in a distributed computing cluster

Inventor: Karthik Kambatla (Mountain View, CA)
Assignee: CLOUDERA, INC.
G06F9/4881G06F9/50G06F9/5038G06F9/505G06F9/5083G06F2209/483G06F2209/5021
View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 12,223,349
App. No.
17/379,742
Granted
Feb 11, 2025
Kind
B2
Abstract

Embodiments are disclosed for a utilization-aware approach to cluster scheduling, to address this resource fragmentation and to improve cluster utilization and job throughput. In some embodiments a resource manager at a master node considers actual usage of running tasks and schedules opportunistic work on underutilized worker nodes. The resource manager monitors resource usage on these nodes and preempts opportunistic containers in the event this over-subscription becomes untenable. In doing so, the resource manager effectively utilizes wasted resources, while minimizing adverse effects on regularly scheduled tasks.

Claims (42)

1. A method comprising:

receiving information indicative of actual computing resource utilization at a worker node, the worker node having allocated a first-tier resource container to process a first task and a second-tier resource container to process a second task;

allocating an opportunistic third-tier resource container at the worker node to process a third task, in response to determining that the actual computing resource utilization at the worker node has decreased below a first threshold since a time that the first-tier resource container and the second-tier resource container were allocated based on summing (i) unutilized resource slack occurring within at least one of the first-tier resource container and the second-tier resource container and (ii) an amount of unutilized unallocated resources at the worker node;

wherein the opportunistic third-tier resource container includes at least a portion of the unutilized resource slack occurring within the at least one of the first-tier resource container and the second-tier resource container, and wherein the opportunistic third-tier resource container, having a lower tier status than the first-tier resource container and the second-tier resource container, is subject to de-allocation to the first-tier resource container and the second-tier resource container if the actual computing resource utilization at the worker node rises above a second threshold that is offset below a utilization capacity associated with the worker node;

determining whether the actual computing resource utilization is between the second threshold and the utilization capacity associated with the worker node; and

deallocating the opportunistic third-tier resource container in response to determining that the actual computing resource utilization is between the second threshold and the utilization capacity such that said deallocating occurs prior to the actual computing resource utilization reaching the utilization capacity associated with the worker node.

2. The method of claim 1 , wherein the opportunistic third-tier resource container is allocated to process the third task after determining that the third task allows processing by opportunistic third-tier resources.

3. The method of claim 1 , further comprising:

determining whether to de-allocate the opportunistic third-tier resource container based on a periodic heartbeat signal of the received information.

4. The method of claim 3 , wherein determining whether to de-allocate the opportunistic third-tier resource container includes:

determining whether the actual computing resource utilization has risen above the second threshold when the periodic heartbeat signal is received.

5. The method of claim 1 , wherein the received information includes values for the first threshold and/or the second threshold.

6. The method of claim 1 , wherein the first threshold is based on the utilization capacity associated with the worker node.

7. The method of claim 6 , wherein the first threshold is further based on a variable over-allocation parameter.

8. The method of claim 1 , wherein the second threshold is further based on a variable preemption parameter.

9. The method of claim 1 , further comprising:

setting the first threshold and/or the second threshold based on the received information.

10. The method of claim 1 , further comprising:

adjusting the first threshold and/or the second threshold in response to changes in the actual computing resource utilization.

11. The method of claim 1 , wherein the second threshold is higher than the first threshold.

12. The method of claim 1 , wherein the first threshold and/or the second threshold are specific to the worker node and are different from thresholds of other worker nodes in a distributed computing cluster.

13. A system comprising:

a processor; and

a memory having instructions stored thereon, which when executed by the processor, cause the system to:

receive information indicative of actual computing resource utilization at a worker node, the worker node having allocated a first-tier resource container to process a first task and a second-tier resource container to process a second task; and

allocate an opportunistic third-tier resource container at the worker node to process a third task, in response to determining that the actual computing resource utilization at the worker node has decreased below a first threshold since a time that the first-tier resource container and the second-tier resource container were allocated based on summing (i) unutilized resource slack occurring within at least one of the first-tier resource container and the second-tier resource container and (ii) an amount of unutilized unallocated resources at the worker node;

wherein the opportunistic third-tier resource container includes at least a portion of the unutilized resource slack occurring within the at least one of the first-tier resource container and the second-tier resource container, and wherein the opportunistic third-tier resource container, having a lower tier status than the first-tier resource container and the second-tier resource container, is subject to de-allocation to the first-tier resource container and the second-tier resource container if the actual computing resource utilization at the worker node rises above a second threshold that is offset below a utilization capacity associated with the worker node;

determine whether the actual computing resource utilization is between the second threshold and the utilization capacity associated with the worker node; and

deallocate the opportunistic third-tier resource container in response to determining that the actual computing resource utilization is between the second threshold and the utilization capacity such that said deallocation occurs prior to the actual computing resource utilization reaching the utilization capacity associated with the worker node.

14. The system of claim 13 , wherein the memory has further instructions stored thereon, which when executed by the processor, cause the system to further:

allocate a second opportunistic third-tier resource container at the worker node to process a fourth task in response to determining that the actual computing resource utilization at the worker node is below a third threshold,

wherein the second opportunistic third-tier resource container includes at least a second portion of the unutilized resource slack occurring within the at least one of the first-tier resource container and the second-tier resource container, and wherein the second opportunistic third-tier resource container, having a lower tier status than the first-tier resource container and/or the second-tier resource container, is subject to de-allocation if the actual computing resource utilization at the worker node rises above a fourth threshold.

15. The system of claim 14 , wherein the third threshold is the same as the first threshold, and wherein the fourth threshold is the same as the second threshold.

16. The system of claim 13 , wherein computing resources at the worker node include one or more of processing, memory, data storage, I/O, and/or network resources.

17. A non-transitory computer readable medium storing instructions, execution of which by a computer system, cause the computer system to:

receive information indicative of actual computing resource utilization at a worker node, the worker node having allocated a first-tier resource container to process a first task and a second-tier resource container to process a second task; and

allocate an opportunistic third-tier resource container at the worker node to process a third task in response to determining that the actual computing resource utilization at the worker node has decreased below a first threshold since a time that the first-tier resource container and the second-tier resource container were allocated based on summing (i) unutilized resource slack occurring within at least one of the first-tier resource container and the second-tier resource container and (ii) an amount of unutilized unallocated resources at the worker node;

wherein the opportunistic third-tier resource container includes at least a portion of the unutilized resource slack occurring within the at least one of the first-tier resource container and the second-tier resource container, and wherein the opportunistic third-tier resource container, having a lower tier status than the first-tier resource_container and the second-tier resource container, is subject to de-allocation to the first-tier resource container and the second-tier resource container if the actual computing resource utilization at the worker node rises above a second threshold that is offset below a utilization capacity associated with the worker node;

determine whether the actual computing resource utilization is between the second threshold and the utilization capacity associated with the worker node; and

deallocate the opportunistic third-tier resource container in response to determining that the actual computing resource utilization is between the second threshold and the utilization capacity such that said deallocating occurs prior to the actual computing resource utilization reaching the utilization capacity associated with the worker node.

18. The non-transitory computer readable medium of claim 17 , storing further instructions, execution of which by the computer system, cause the computer system to further:

receive a request to process a particular task, wherein the request includes an indication of whether to allow processing of the particular task using the opportunistic third-tier resource container.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jul 23, 2021
From: KAMBATLA, KARTHIK
To: CLOUDERA, INC
Reel/Frame 056959/0515 →
Continuity (4)
Continuation 16797996 · Feb 21, 2020
Continuation 15595713 · May 15, 2017
Provisional Application 62394660 · Sep 14, 2016
Related Publication 20210349755A1 · Nov 11, 2021
References Cited (21)
US 8706798B1 · Suchter · 2014 [cited by examiner]
US 9374243B1 · Certain · 2016 [cited by examiner]
US 10572306B2 · Kambatla · 2020 [cited by applicant]
US 11099892B2 · Kambatla · 2021 [cited by applicant]
US 20060140115A1 · Timus et al. · 2006 [cited by applicant]
US 20080082979A1 · Coppinger et al. · 2008 [cited by applicant]
US 20090007125A1 · Barsness · 2009 [cited by examiner]
US 20100262695A1 · Mays · 2010 [cited by examiner]
US 20110055838A1 · Moyes · 2011 [cited by applicant]
US 20120254822A1 · Zheng et al. · 2012 [cited by applicant]
US 20140245298A1 · Zhou et al. · 2014 [cited by applicant]
US 20150026336A1 · Suchter · 2015 [cited by examiner]
US 20150363238A1 · Bai · 2015 [cited by examiner]
US 20170017521A1 · Gupta et al. · 2017 [cited by applicant]
US 20180074855A1 · Kambatla · 2018 [cited by applicant]
US 20200192703A1 · Kambatla · 2020 [cited by applicant]
International Search Report and Written Opinion of International Application No. PCT/US2017/043137; Date of Mailing: Oct. 26, 2017; 17 pages. [cited by applicant]
Extended European Search Report, EP Patent Application No. 17851228.1 based on PCT/US2017/043137, mailed Apr. 30, 2020, 11 pages. [cited by applicant]
Yao Yi et al.: “Opera: Opportunistic and Efficient Resource Allocation in Hadoop Yarn by Harnessing Idle Resources”, 2016 25th International Conference on Computer Communication and Networks (ICCCN), IEEE, Aug. 1, 2016,… [cited by applicant]
Wei Zhang et al.: “Minimizing Interference and Maximizing Progress for Hadoop Virtual Machines”, ACM Sigmetrics Performance Evaluation Review, Association for Computing Machinery, New York, NY, US, vol. 42, No. 4, Jun. … [cited by applicant]
Karanasos et al.: “Mercury: Hybrid Centralized and Distributed Scheduling in Large Shared Clusters”, Microsoft Corporation, Usenix, the Advanced Computing Systems Association, Annual Technical Conference, Jun. 27, 2015,… [cited by applicant]