IP Library Granted Patent US 9,325,593
Granted Patent B2
US 9,325,593 · App. 14/467,629 · Granted Apr 26, 2016

Systems, methods, and devices for dynamic resource monitoring and allocation in a cluster system

Inventors: Sean Andrew Suchter (Los Altos Hills, CA); Charles C. Carson, Jr. (Cupertino, CA); Kimoon Kim (Cupertino, CA); Choongsoon Chang (Palo Alto, CA); Scott Alexander Banachowski (Mountain View, CA); Judith A. Hay (Basel, CH)
Assignee: Pepperdata, Inc.
H04L43/08G06F9/5038G06F9/5066H04L41/24H04L43/0876H04L47/70G06F2209/508
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 9,325,593
App. No.
14/467,629
Granted
Apr 26, 2016
Kind
B2
Abstract

In an embodiment, the systems, methods, and devices disclosed herein comprise a computer resource monitoring and allocation system. In an embodiment, the resource monitoring and allocation system can be configured to allocate computer resources that are available on various nodes of a cluster to specific jobs and/or sub-jobs and/or tasks and/or processes.

Claims (50)

1. A hadoop computer cluster comprising:

one or more processors of a master node, wherein the master node comprises a supervisor controller;

one or more processors of a plurality of computing system nodes, the one or more processors of the plurality of computing system nodes configured to perform computing processes on received sub-jobs, wherein each computing system node comprises an agent controller;

wherein each agent controller comprises:

a monitoring interface configured to monitor system resources utilization by sub-jobs of its respective computing system node; and

a reporting controller configured to transmit the monitored system resources utilization to the supervisor controller in substantially real-time;

wherein the supervisor controller is configured to assign an additional sub-job to a first computing system node based on determining that the utilization of a first electronic random access memory capacity of the first computing system node is below a threshold level, the determining based on the monitored system resources utilization transmitted from the reporting controller of the first computing system node to the supervisor controller,

wherein the supervisor controller is configured to monitor a second electronic random access memory capacity of a second computing system node,

wherein the supervisor controller is configured to prevent assignment of additional sub-jobs to a second computing system node based on determining that utilization of the second electronic random access memory capacity is at or above a threshold value, and

wherein the master node and the plurality of computing system nodes include a computer processor and an electronic storage medium.

2. The hadoop computer cluster of claim 1 , further comprising a management controller, wherein the management controller comprises a job tracker or a yarn resource manager, wherein the first computer system node comprises a task tracker or a yarn node manager configured to communicate data to the job tracker or the yarn resource manager, the communicated data comprising a number of available slots or containers for processing tasks on the first computer system node.

3. The hadoop computer cluster of claim 2 , wherein the supervisor controller is configured to assign the additional sub-job to the first computing system node based on determining that the utilization of at least one system resource of the first computing system node is below a threshold level and based on determining that the number of available slots or containers is zero on the first computing system node.

4. The hadoop computer cluster of claim 1 , wherein the hadoop computer cluster operates in a non-virtual environment.

5. The hadoop computer cluster of claim 1 , wherein the hadoop computer cluster operates in a virtual environment.

6. A supervisor controller configured to dynamically manage assignment of job processes in a hadoop computer cluster, the supervisor controller comprising:

a management controller interface configured to communicate with a management controller to access data representing an assignment of a plurality of job processes across a plurality of computer system nodes in the computer cluster;

an agent controller interface configured to communicate with an agent controller operating on a first computing system node to receive data representing utilization of system resources on the first computing system node, wherein the agent controller interface is further configured to communicate with an agent controller operating on a second computing system node to receive data representing utilization of system resources on the second computing system node; and

a system resource allocation engine configured to assign an additional job process to the first computing system node based on determining that the utilization of a first electronic random access memory capacity of the first computing system node is below a threshold level; and

wherein one or more computer processors and one or more electronic storage medium are configured to operate the supervisor controller;

wherein the system resource allocation engine is configured to monitor a second electronic random access memory capacity of a second computing system node, and

wherein the system resource allocation engine is configured to prevent assignment of additional job processes to a second computing system node based on determining that utilization of the second electronic random access memory capacity is at or above a threshold level.

7. The supervisor controller of claim 6 , wherein the management controller comprises a job tracker or a yarn resource manager, wherein the plurality of computer system nodes comprise a task tracker or a yarn node manager configured to communicate data to the job tracker or the yarn resource manager, the communicated data comprising a number of available slots or containers for processing tasks on a particular computer system node.

8. The supervisor controller of claim 7 , wherein the supervisor controller is configured to assign the additional job process to the first computing system node based on determining that the utilization of at least one system resource of the computing system node is below a threshold level and based on determining that the number of available slots or containers is zero for the first computing system node.

9. The supervisor controller of claim 6 , wherein the supervisor controller is configured to assign the additional job process to the first computing system node based on determining that the utilization of a first system resource of the first computing system node is below a first threshold level and based on determining that the utilization of a second system resource of the first computing system node is above a second threshold level, the determining based on the data representing utilization of system resources received from the agent controller by the supervisor controller.

10. The supervisor controller of claim 9 , wherein the additional job process assigned by the supervisor controller requires greater utilization of the first system resource as compared to the utilization of the second system resource.

11. A hadoop computer cluster comprising:

one or more processors of a master node, wherein the master node comprises a supervisor controller;

one or more processors of a plurality of computing system nodes, the one or more processors of the plurality of computing system nodes configured to perform computing processes on received sub-jobs, wherein each computing system node comprises an agent controller;

wherein each agent controller comprises:

a monitoring interface configured to monitor system resources utilization by sub-jobs of its respective computing system node; and

a reporting controller configured to transmit the monitored system resources utilization to the supervisor controller in substantially real-time;

wherein the supervisor controller is configured to assign an additional sub-job to a first computing system node based on determining that the utilization of a first CPU processor capacity of the first computing system node is below a threshold level, the determining based on the monitored system resources utilization transmitted from the reporting controller of the first computing system node to the supervisor controller;

wherein the supervisor controller is configured to monitor a second CPU processor capacity of a second computing system node,

wherein the supervisor controller is configured to prevent assignment of additional sub-jobs to the second computing system node based on determining that utilization of the second CPU processor capacity is at or above a threshold value, and

wherein the master node and the plurality of computing system nodes include a computer processor and an electronic storage medium.

12. The hadoop computer cluster of claim 11 further comprising a management controller, wherein the management controller comprises a job tracker or a yarn resource manager, wherein the plurality of computing system nodes comprise task trackers or yarn node managers configured to communicate data to the job tracker or the yarn resource manager, the communicated data comprising a number of available slots or containers for processing tasks on a particular computer system node.

13. The hadoop computer cluster of claim 12 , wherein the supervisor controller is configured to assign the additional sub-job to the first computing system node based on determining that the utilization of the at least one system resource of the computing system node is below the threshold level and based on determining that the number of available slots or containers is zero on the first computing system node.

14. The hadoop computer cluster of claim 11 , wherein the hadoop computer cluster operates in a non-virtual environment.

15. The hadoop computer cluster of claim 11 , wherein the hadoop computer cluster operates in a virtual environment.

16. A supervisor controller configured to manage assignment of job processes in a computer cluster, the supervisor controller comprising:

a management controller interface configured to communicate with a management controller to access data representing an assignment of a plurality of job processes across a plurality of computing system nodes in the computer cluster;

an agent controller interface configured to communicate with an agent controller operating on a first computing system node to receive data representing utilization of system resources by the plurality of job processes operating on the first computing system node, wherein the agent controller interface is further configured to communicate with an agent controller operating on a second computing system node to receive data representing utilization of system resources on the second computing system node; and

a system resource allocation engine configured to assign an additional job process to a first computing system node based on determining that the utilization of a first CPU processor capacity of the first computing system node is below a threshold level; and

wherein one or more computer processors and one or more electronic storage medium are configured to operate the supervisor controller;

wherein the system resource allocation engine is configured to monitor a second CPU processor capacity of a second computing system node, and

wherein the system resource allocation engine is configured to prevent assignment of additional job processes to the second computing system node based on determining that utilization of the second CPU processor capacity is at or above a threshold value.

17. The supervisor controller of claim 16 , wherein the management controller comprises a job tracker or a yarn resource manager, wherein the plurality of computer system nodes comprise task trackers or yarn node managers configured to communicate data to the job tracker or the yarn resource manager, the communicated data comprising a number of available slots or containers for processing tasks on a particular computer system node.

18. The supervisor controller of claim 17 , wherein the supervisor controller is configured to assign the additional job process to the first computing system node based on determining that the utilization of the at least one system resource of the computing system node is below the threshold level and based on determining that the number of available slots or containers is zero for the computing system node.

19. The supervisor controller of claim 16 , wherein the supervisor controller is configured to assign the additional job process to the first computing system node based on determining that the utilization of a first system resource of the first computing system node is below a first threshold level and based on determining that the utilization of a second system resource of the first computing system node is above a second threshold level, the determining based on the data representing utilization of system resources received from the agent controller by the supervisor controller.

20. The supervisor controller of claim 19 , wherein the additional job process assigned by the supervisor controller requires greater utilization of the first system resource as compared to the utilization of the second system resource.

Assignments (3)
CHANGE OF ASSIGNEE ADDRESS Recorded Nov 1, 2024
From: PEPPERDATA, INC.
To: PEPPERDATA, INC.
Reel/Frame 069290/0110 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 25, 2015
From: SUCHTER, SEAN ANDREW; CARSON, CHARLES C., JR.; KIM, KIMOON; CHANG, CHOONGSOON; BANACHOWSKI, SCOTT ALEXANDER; HAY, JUDITH A.
To: PEPPERDATA, INC.
Reel/Frame 036418/0198 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 27, 2014
From: SUCHTER, SEAN ANDREW; CARSON, CHARLES C., JR.; KIM, KIMOON; CHANG, CHOONGSOON; BANACHOWSKI, SCOTT ALEXANDER; HAY, JUDITH A.
To: PEPPERDATA, INC.
Reel/Frame 033623/0408 →
Continuity (9)
Continuation 14194406 · Feb 28, 2014
Continuation 14053044 · Oct 14, 2013
Provisional Application 61841007 · Jun 28, 2013
Provisional Application 61841074 · Jun 28, 2013
Provisional Application 61841127 · Jun 28, 2013
Provisional Application 61841025 · Jun 28, 2013
Provisional Application 61841106 · Jun 28, 2013
Provisional Application 61841061 · Jun 28, 2013
Related Publication 20150026336A1 · Jan 22, 2015