US2015186189A1PendingUtilityA1
Managing array computations during programmatic run-time in a distributed computing environment
Assignee: HEWTETT PACKARD DEV COMPANY L PPriority: Jul 30, 2012Filed: Mar 10, 2015Published: Jul 2, 2015
Est. expiryJul 30, 2032(~6 yrs left)· nominal 20-yr term from priority
Inventors:Shivaram VenkataramanIndrajit RoyMehul A. ShahRobert SchreiberNathan Lorenzo BinkertParthasarathy Ranganathan
G06F 9/4881G06F 9/5066G06F 9/5088
45
PatentIndex Score
0
Cited by
0
References
0
Claims
Abstract
A plurality of array partitions are defined for use by a set of tasks of the program run-time. The array partitions can be determined from one or more arrays that are utilized by the program at run-time. Each of the plurality of computing devices are assigned to perform one or more tasks in the set of tasks. By assigning each of the plurality of computing devices to perform one or more tasks, an objective to reduce data transfer amongst the plurality of computing devices can be implemented.
Claims
exact text as granted — not AI-modified1 - 15 . (canceled)
16 . A method for executing a program comprising a set of tasks, the method performed by one or more processors of a master node and comprising:
assigning a plurality of the set of tasks to each respective one of a plurality of worker nodes; sequencing the plurality of tasks for execution by the respective worker node by (i) identifying a first task in the plurality of tasks that yields an array, (ii) identifying a second task in the plurality of tasks that utilizes the array, and (iii) assigning the first task to be performed by the respective worker node before the second task.
17 . The method of claim 16 , further comprising:
identifying an update to the array; and implementing a callback on the respective node to perform the second task utilizing the updated array.
18 . The method of claim 16 , further comprising:
maintaining, in a memory resource, metadata defining a current execution state of the program.
19 . The method of claim 18 , further comprising:
backing up the respective worker node after completion of each of the plurality of tasks.
20 . The method of claim 18 , wherein the metadata further defines an execution state for each of the plurality of worker nodes.
21 . The method of claim 20 , further comprising:
transmitting a periodic message to the plurality of worker nodes to determine the current execution state for each of the plurality of worker nodes.
22 . The method of claim 21 , further comprising:
determining, based on the periodic message, that the respective worker node has failed; in response to determining that the respective node has failed, identifying, from the metadata, a previous execution state for the respective worker node; and restarting the respective worker node based on the previous execution state.
23 . A master node comprising:
at least one processor; and at least one memory resource storing instructions for executing a program comprising a set of tasks, wherein the instructions, when executed by the at least one processor, cause the master node to:
assign a plurality of the set of tasks to each respective one of a plurality of worker nodes;
sequence the plurality of tasks for execution by the respective worker node by (i) identifying a first task in the plurality of tasks that yields an array, (ii) identifying a second task in the plurality of tasks that utilizes the array, and (iii) assigning the first task to be performed by the respective worker node before the second task.
24 . The master node of claim 23 , wherein the instructions, when executed by the at least one processor, further cause the master node to:
identify an update to the array; and implement a callback on the respective node to perform the second task utilizing the updated array.
25 . The master node of claim 23 , wherein the instructions, when executed by the at least one processor, further cause the master node to:
maintain, in the at least one memory resource, metadata defining a current execution state of the program.
26 . The master node of claim 25 , wherein the instructions, when executed by the at least one processor, further cause the master node to:
back up the respective worker node after completion of each of the plurality of tasks.
27 . The master node of claim 25 , wherein the metadata further defines an execution state for each of the plurality of worker nodes.
28 . The master node of claim 27 , wherein the instructions, when executed by the at least one processor, further cause the master node to:
transmit a periodic message to the plurality of worker nodes to determine the current execution state for each of the plurality of worker nodes.
29 . The master node of claim 28 , wherein the instructions, when executed by the at least one processor, further cause the master node to:
determine, based on the periodic message, that the respective worker node has failed; in response to determining that the respective node has failed, identify, from the metadata, a previous execution state for the respective worker node; and restart the respective worker node based on the previous execution state.
30 . A non-transitory computer readable medium storing instructions for executing a program comprising a set of tasks, wherein the instructions, when executed by at least one processor of a master node, cause the master node to:
assign a plurality of the set of tasks to each respective one of a plurality of worker nodes; sequence the plurality of tasks for execution by the respective worker node by (i) identifying a first task in the plurality of tasks that yields an array, (ii) identifying a second task in the plurality of tasks that utilizes the array, and (iii) assigning the first task to be performed by the respective worker node before the second task.
31 . The non-transitory computer readable medium of claim 30 , wherein the instructions, when executed by at least one processor of a master node, further cause the master node to:
identify an update to the array; and implement a callback on the respective node to perform the second task utilizing the updated array.
32 . The non-transitory computer readable medium of claim 30 , wherein the instructions, when executed by at least one processor of a master node, further cause the master node to:
maintain, in a memory resource, metadata defining a current execution state of the program and an execution state for each of the plurality of worker nodes.
33 . The non-transitory computer readable medium of claim 32 , wherein the instructions, when executed by at least one processor of a master node, further cause the master node to:
back up the respective worker node after completion of each of the plurality of tasks.
34 . The non-transitory computer readable medium of claim 33 , wherein the instructions, when executed by at least one processor of a master node, further cause the master node to:
transmit a periodic message to the plurality of worker nodes to determine the current execution state for each of the plurality of worker nodes.
35 . The non-transitory computer readable medium of claim 34 , wherein the instructions, when executed by at least one processor of a master node, further cause the master node to:
determine, based on the periodic message, that the respective worker node has failed; in response to determining that the respective node has failed, identify, from the metadata, a previous execution state for the respective worker node; and restart the respective worker node based on the previous execution state.Join the waitlist — get patent alerts
Track US2015186189A1 — get alerts on status changes and closely related new filings.
We store only your email — no account needed. See our privacy policy.