IP Library › Patent Application 18568050
Patent Application
App. No. 18/568,050

METHOD AND SYSTEM FOR DISTRIBUTED WORKLOAD PROCESSING

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 None
App. No.
18/568,050
Abstract

A method and system for distributing a compute model and data to process to heterogeneous and distributed compute devices. The compute model and a portion of the data is processed on a benchmark system and the timing used to make a job execution speed estimate for each compute device. Compute devices are selected and assigned data chunks based on the estimate so distributed processing is completed within a predefined time period. The compute model and data chunks can be sent to the respective compute devices using separate processes, such as a payload manager configured to transfer compute jobs to remote devices and a messaging engine configured to transfer data messages, and where the payload manager and messaging engine communicate with corresponding software engines on the compute devices.

Claims (50)

1 - 19 . (canceled)

20 . A method for distributing data processing jobs to a plurality of compute devices in a distributed computing environment, the method performed by a computer platform comprising at least one computer system, the method comprising the steps of:

receiving at the computer platform a job payload comprising a compute model and job data to be processed by the compute model;

dividing the job data into a plurality of data chunks;

sending the compute model and a benchmark data set extracted from the job data to a benchmark compute platform;

receiving from the benchmark compute platform an indication of benchmark job computation speed for processing of the benchmark data by the compute model;

selecting a cluster of compute devices from the plurality of compute devices and allocating each of the plurality of data chunks to a respective compute device in the cluster such that for each respective compute device in the cluster, an estimated duration to complete processing of allocated data chunks with the compute model is less than a first predefined duration, and wherein the estimated duration is determined based on the benchmark job computation speed;

sending each respective compute device in the cluster the compute model and the data chunks allocated to the respective compute device in the cluster; and

receiving from respective compute devices in the cluster compute result data sets, each compute result data set associated with a respective data chunk.

21 . The method of claim 20 , wherein the selection of the cluster of compute devices and allocation of the data chunks to the respective compute devices in the cluster is such that for each respective compute device in the cluster, the estimated duration to complete processing of allocated data chunks with the compute model is greater than a second predefined duration, the second predefined duration being less than the first predefined duration.

22 . The method of claim 20 , further comprising the steps of:

identifying an insufficient compute resource condition wherein a cluster of compute devices cannot be selected from available compute devices in the plurality of compute devices and data chunks to be allocated to the compute devices in the cluster so as to provide an estimated duration to complete processing of the data chunks that is less than the first predefined duration; and

in response to the insufficient compute resource condition, creating a virtual compute device and making the virtual compute device available for inclusion in the cluster of compute devices.

23 . The method of claim 20 , wherein the step of sending each respective compute device in the cluster the compute model and the data chunks allocated to that respective compute device comprises the steps of:

sending the compute model to the respective compute devices using a first data transfer process; and

sending the data chunks to the respective compute devices using a second data transfer process different from the first data transfer process.

24 . The method of claim 20 ,

the benchmark platform comprises a plurality of benchmark devices each having a respective device type including at least a first device of a first device type and a second device of a second device type different from the first device type; the indication of benchmark job computation speed comprising a plurality of benchmark speeds, each benchmark speed associated with a respective benchmark device;

the estimated duration to complete processing of allocated data chunks by the compute model for a respective compute device in the cluster being determined based on the benchmark speed of a benchmark device in the plurality of benchmark devices that is most similar relative to remaining benchmark devices to the respective compute device in the cluster as indicated by a similarity function.

25 . The method of claim 20 , further comprising the steps of reassigning to a replacement compute device the data chunks allocated to a first compute device in the cluster and for which compute results have not been received by a predetermined deadline.

26 . The method of claim 25 , wherein the replacement compute device is selected from the compute devices in the plurality of compute devices which are not in the cluster, the step of reassigning to the replacement compute device further comprising sending the replacement compute device the compute model.

27 . The method of claim 20 , further comprising the step of:

in response to a health alert condition for a first compute device in the cluster reassigning the data chunks allocated to the first compute device and for which compute results have not been received to a replacement compute device selected from the plurality of compute devices.

28 . The method of claim 27 , further comprising the step of periodically requesting from each respective compute devices in the cluster a response message indicating health status of the respective compute device.

29 . The method of claim 27 , wherein the health alert condition for a respective compute device in the cluster comprises a lack of available resources at the compute device to complete processing of the allocated chunks with the compute model within the first predefined duration.

30 . The method of claim 29 , wherein the available resources comprises at least one of battery charge, available memory, processor temperature, CPU resources, and execution speed of the compute job.

31 . A system for distributing data processing jobs to a plurality of compute devices in a distributed computing environment, the system comprising:

a computer platform comprising at least one programmable computer system having a computer memory storing instructions that are executable to configure the computer platform to:

the memory having computer instructions stored therein which, when executed, configure the processor to:

receive a job payload comprising a compute model and job data to be processed by the compute model;

divide the job data into a plurality of data chunks;

send the compute model and a benchmark data set extracted from the job data to a benchmark compute platform;

receive from the benchmark compute platform an indication of benchmark job computation speed for processing of the benchmark data by the compute model;

select a cluster of compute devices from the plurality of compute devices and allocate each of the plurality of data chunks to a respective compute device in the cluster such that for each respective compute device in the cluster, an estimated duration to complete processing of allocated data chunks with the compute model is less than a first predefined duration, and wherein the estimated duration is determined based on the benchmark job computation speed;

send to each respective compute device in the cluster the compute model and the data chunks allocated to the respective compute device in the cluster; and

receive from respective compute devices in the cluster compute result data sets, each compute result data set associated with a respective data chunk.

32 . The system of claim 31 , wherein the computer platform is configured to select the cluster of compute devices and allocate the data chunks to the respective compute devices in the cluster is such that for each respective compute device in the cluster, the estimated duration to complete processing of allocated data chunks with the compute model is greater than a second predefined duration, the second predefined duration being less than the first predefined duration.

33 . The system of claim 31 , the computer instructions are executable to further configure the computer platform to:

identify an insufficient compute resource condition wherein a cluster of compute devices cannot be selected from available compute devices in the plurality of compute devices and data chunks to be allocated to the compute devices in the cluster so as to provide an estimated duration to complete processing of the data chunks that is less than the first predefined duration; and

create in response to the insufficient compute resource condition a virtual compute device, which virtual compute device is available for inclusion in the cluster of compute devices.

34 . The system of claim 31 , wherein computer platform is configured send the compute model to the respective compute devices using a first data transfer process and sending the data chunks to the respective compute devices using a second data transfer process different from the first data transfer process.

35 . The system of claim 31 , wherein the benchmark platform comprises a plurality of benchmark devices each having a respective device type including at least a first device of a first device type and a second device of a second device type different from the first device type; the indication of benchmark job computation speed comprising a plurality of benchmark speeds, each benchmark speed associated with a respective benchmark device; and

the estimated duration to complete processing of allocated data chunks by the compute model for a respective compute device in the cluster is determined based on the benchmark speed of a benchmark device in the plurality of benchmark devices that is most similar relative to remaining benchmark devices to the respective compute device in the cluster as indicated by a similarity function.

36 . The system of claim 31 , wherein the computer instructions are executable to further configure the computer platform to reassign to a replacement compute device the data chunks allocated to a first compute device in the cluster and for which compute results have not been received by a predetermined deadline.

37 . The system of claim 36 , wherein the computer instructions configure the computer platform to select the replacement compute device from the compute devices in the plurality of compute devices which are not in the cluster and send the replacement compute device the compute model.

38 . The system of claim 31 , wherein the computer instructions are executable to further configure the computer platform to:

in response to receipt of a health alert condition for a first compute device in the cluster, reassign the data chunks allocated to the first compute device and for which compute results have not been received to a replacement compute device selected from the plurality of compute devices.

39 . The system of claim 38 , wherein the computer instructions are executable to further configure the computer platform to periodically request from each respective compute device in the cluster a response message indicating health status of the respective compute device.

40 . The system of claim 38 , wherein the health alert condition for a respective compute device in the cluster comprises a lack of available resources at the compute device to complete processing of the allocated chunks with the compute model within the first predefined duration.

41 . The system of claim 40 , wherein the available resources comprises at least one of battery charge, available memory, processor temperature, CPU resources, and execution speed of the compute job.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 22, 2024
From: GROSS, TYLER; FELICE, RONALD A.
To: SAILION INC.
Reel/Frame 067494/0523 →