IP Library › Granted Patent US 9,244,751
Granted Patent B2
US 9,244,751 · App. 14/118,520 · Granted Jan 26, 2016

Estimating a performance parameter of a job having map and reduce tasks after a failure

Inventors: Ludmila Cherkasova (Sunnyvale, CA); Abhishek Verma (Champaign, IL)
Assignee: Hewlett Packard Enterprise Development LP
G06F11/0715G06F11/3404G06F11/3447G06F11/2041G06F11/3419G06F11/3433G06F2201/865
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,244,751
App. No.
14/118,520
Filed
Nov 18, 2013
Granted
Jan 26, 2016
Kind
B2
Art Unit
2113
USPC
714/4.12
Abstract

A job profile includes characteristics of a job to be executed, where the characteristics of the job profile relate to map tasks and reduce tasks of the job, and where the map tasks produce intermediate results based on input data, and the reduce tasks produce an output based on the intermediate results. In response to a failure in a system, numbers of failed map tasks and reduce tasks of the job based on a time of the failure are computed, and numbers of remaining map tasks and reduce tasks are computed. A performance model is provided, and a performance parameter of the job is estimated using the performance model.

Claims (44)

1. A method comprising:

receiving, by a system comprising a processor, a job profile that includes characteristics of a job to be executed, wherein the characteristics of the job profile relate to map tasks and reduce tasks of the job, wherein the map tasks produce intermediate results based on input data, and the reduce tasks produce an output based on the intermediate results;

in response to an indication of a failure in the system associated with execution of the job, computing, by the system, numbers of failed map tasks and failed reduce tasks of the job based on a time of the failure, wherein the indication of the failure comprises an indication that a computing node of plural computing nodes has failed;

computing, by the system, numbers of remaining map tasks and remaining reduce tasks based on the numbers of failed map tasks and failed reduce tasks;

providing, by the system, a performance model based on the job profile, the computed numbers of remaining map tasks and remaining reduce tasks, and an allocated amount of resources for the job, wherein providing the performance model is based on the allocated amount of resources that has been reduced from a previous allocation of resources due to the failure; and

estimating, by the system, a performance parameter of the job using the performance model.

2. The method of claim 1 , wherein the indication comprises a time indication regarding the time of the failure.

3. The method of claim 2 , further comprising:

based on the time of failure, determining whether the failure is during a map stage or reduce stage, where the map tasks of the job are performed during the map stage, and wherein the reduce tasks of the job are performed in the reduce stage.

4. The method of claim 3 , wherein, based on the failure occurring during the map stage, the computed number of failed reduce tasks is zero.

5. The method of claim 3 , wherein, based on the failure occurring during the reduce stage, the computed number of failed map tasks is equal to a number of the map tasks of the job in the map stage divided by a number of worker nodes for executing map and reduce tasks.

6. The method of claim 1 , wherein estimating the performance parameter comprises estimating a time duration of the job.

7. An article comprising at least one non-transitory machine-readable storage medium storing instructions that upon execution cause a system to:

receive a job profile that includes characteristics of a job to be executed, wherein the characteristics of the job profile relate to map tasks and reduce tasks of the job, wherein the map tasks produce intermediate results based on input data, and the reduce tasks produce an output based on the intermediate results;

in response to an indication of a failure in the system associated with execution of the job, compute numbers of failed map tasks and failed reduce tasks of the job based on a time of the failure;

compute numbers of remaining map tasks and remaining reduce tasks based on the numbers of failed map tasks and failed reduce tasks;

determine whether resources of the system are to be replenished after the failure;

provide a performance model based on the job profile, the computed numbers of remaining map tasks and remaining reduce tasks, and an allocated amount of resources for the job, wherein providing the performance model is based on the allocated amount of resources after the replenishing of resources; and

estimate a performance parameter of the job using the performance model.

8. The article of claim 7 , wherein the indication of the failure comprises an indication that a computing node of plural computing nodes has failed.

9. The article of claim 7 , wherein the indication comprises a time indication regarding the time of the failure.

10. The article of claim 7 , wherein estimating the performance parameter comprises estimating a time duration of the job.

11. A system comprising:

a storage medium to store a job profile that includes characteristics of a job to be executed, wherein the characteristics of the job profile relate to map tasks and reduce tasks of the job, wherein the map tasks produce intermediate results based on input data, and the reduce tasks produce an output based on the intermediate results; and

at least one processor to:

detect occurrence of a failure in the system;

in response to detecting the failure, determine whether the failure occurred during a map stage or a reduce stage, wherein the map tasks are to be performed during the map stage, and the reduce tasks are to be performed during the reduce stage;

compute a number of failed map tasks and a number of failed reduce tasks dependent upon whether the failure occurred during the map stage or the reduce stage;

compute a number of remaining map tasks and a number of remaining reduce tasks based on the numbers of failed map tasks and failed reduce tasks;

update a performance model using the numbers of remaining map tasks and remaining reduce tasks; and

estimate a performance parameter of the job using the updated performance model.

12. The system of claim 11 , wherein the at least one processor is to further: receive an indication of a time of the failure,

wherein determining whether the failure occurred during the map stage or the reduce stage is based on comparing the time of the failure with a time bound associated with the map stage.

13. The system of claim 11 , wherein the at least one processor is to maintain, in the performance model, an allocation of resources equal to an initial allocation of resources for the job in response to an indication that failed resources are to be replenished.

14. The system of claim 11 , wherein the performance parameter is a time duration of the job.

15. A system comprising:

a storage medium to store a job profile that includes characteristics of a job to be executed, wherein the characteristics of the job profile relate to map tasks and reduce tasks of the job, wherein the map tasks produce intermediate results based on input data, and the reduce tasks produce an output based on the intermediate results; and

at least one processor to:

detect occurrence of a failure in the system;

in response to detecting the failure, compute a number of failed map tasks and a number of failed reduce tasks based on a time of the failure;

compute a number of remaining map tasks and a number of remaining reduce tasks based on the numbers of failed map tasks and failed reduce tasks;

update a performance model using the numbers of remaining map tasks and remaining reduce tasks and in response to an indication that failed resources are not to be replenished, wherein the update of the performance model includes reducing allocations of resources; and

estimate a performance parameter of the job using the updated performance model.

16. The system of claim 15 , wherein the performance parameter is a time duration of the job.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 9, 2015
From: HEWLETT-PACKARD DEVELOPMENT COMPANY, L.P.
To: HEWLETT PACKARD ENTERPRISE DEVELOPMENT LP
Reel/Frame 037079/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 23, 2014
From: CHERKASOVA, LUDMILA; VERMA, ABHISHEK
To: HEWLETT-PACKARD DEVELOPMENT COMPANY, L.P.
Reel/Frame 032032/0576 →
Continuity (1)
Related Publication 20140089727A1 · Mar 27, 2014