IP Library Granted Patent US 12,210,988
Granted Patent B1
US 12,210,988 · App. 17/681,230 · Granted Jan 28, 2025

Opportunistic job processing using worker processes comprising instances of executable processes created by work order binary code

Inventors: David Konerding (San Mateo, CA); Jordan M. Breckenridge (Menlo Park, CA); Daniel Belov (London, GB)
Assignee: Google LLC
G06Q10/06311G06F9/452G06F9/468G06F9/4843G06F9/485G06F9/5072G06F9/541G06Q10/06315
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,210,988
App. No.
17/681,230
Granted
Jan 28, 2025
Kind
B1
Abstract

A global-level manager access a work order from a client and parameters associated with the work order. A service level agreement to meet the work order parameters is determined. The service level agreement includes a price. An indication is received from the client that the service level agreement is accepted. The one or more input files are partitioned into multiple shards, and the work order into multiple jobs. The jobs are distributed among a plurality of clusters to be processed using underutilized computing resources in the clusters. The job outputs are combined to form the work order output. The jobs are monitored to insure that the deadline for completion of the work order will be met.

Claims (52)

1. A method for processing data, comprising:

receiving, at a computing system, a work order submitted by a client device, the work order comprising an input file and a binary code;

generating, based on the input file and binary code, one or more jobs;

partitioning the input file into one or more input shards;

processing, using processes running on a group of computers in the computing system, the one or more jobs and input shards, and

wherein processing comprises:

assigning, by a first process, the one or more jobs to a second process, the first process being a higher priority process than the second process,

instantiating, by the second process, one or more native clients, each native client hosting a worker process to perform an assigned work assignment associated with the one or more jobs, each worker process being associated with a given input shard and comprising an instance of at least one executable process created by the binary code,

outputting, by each worker process, one or more output shards to the second process,

providing, by the second process, the one or more output shards to a storage element, and

providing the one or more output shards from the storage element to the client device.

2. The method of claim 1 , wherein the first process comprises a cluster-level manager, the second process comprises a task-level manager and the third process comprises a global-level manager.

3. The method of claim 1 , wherein instantiating comprises creating the one or more native client based on a number of jobs generated.

4. The method of claim 1 , wherein partitioning comprises partitioning the input file into the one or more input shards based on different data types.

5. The method of claim 4 , wherein the different data types are dependent on data size.

6. The method of claim 4 , wherein each different data type is associated with a given one of the prefix values.

7. A system comprising:

one or more processing devices; and

one or more memories storing instructions that, when executed by the one or more processing devices cause the one or more processing devices to:

receive a work order submitted by a client device, the work order comprising an input file and a binary code;

generate, based on the input file and binary code, one or more jobs;

partition the input file into one or more input shards;

process, using processes running on a group of computers that form a computing system, the one or more jobs and one or more input shards;

wherein to process comprises:

assigning, by a first process, the one or more jobs to a second process, the first process being a higher priority process than the second process,

instantiating, by the second process, one or more native clients, each native client hosting a worker process to perform an assigned work assignment associated with the one or more jobs, each worker process being associated with a given input shard and comprising an instance of at least one executable process created by the binary code,

outputting, by each worker process, one or more output shards to the second process,

providing, by the second process, the one or more output shards to a storage element, and

providing the one or more output shards from the storage element to the client device.

8. The system of claim 7 , wherein the first process comprises a cluster-level manager, the second process comprises a task-level manager and the third process comprises a global-level manager.

9. The system of claim 7 , wherein instantiating comprises creating the one or more native clients based on a number of jobs generated.

10. The system of claim 7 , wherein partitioning comprises partitioning the input file into the one or more input shards based on different types of data.

11. The system of claim 10 , wherein the different types of data is dependent on data size.

12. The system of claim 10 , wherein each different type of data is associated with a given one of the prefix values.

13. A non-transitory computer readable medium storing executable instructions, the executable instructions when executed by one or more processing devices cause the one or more processing devices to perform operations comprising:

receiving, at a computing system, a work order submitted by a client device, the work order comprising an input file and a binary code;

generating, based on the input file and binary code, one or more jobs;

partitioning the input file into one or more input shards;

processing, using processes running on a group of computers in the computing system, the one or more jobs and input shards, and

wherein processing comprises:

assigning, by a first process, the one or more jobs to a second process, the first process being a higher priority process than the second process,

instantiating, by the second process, one or more native clients, each native client hosting a worker process to perform an assigned work assignment associated with the one or more jobs, each worker process being associated with a given input shard and comprising an instance of at least one executable process created by the binary code,

outputting, by each worker process, one or more output shards to the second process,

providing, by the second process, the one or more output shards to a storage element, and

providing the one or more output shards from the storage element to the client device.

14. The non-transitory computer readable medium of claim 13 , wherein the first process comprises a cluster-level manager, the second process comprises a task-level manager and the third process comprises a global-level manager.

15. The non-transitory computer readable medium of claim 13 , wherein instantiating comprises creating the one or more native client based on a number of jobs generated.

16. The non-transitory computer readable medium of claim 13 , wherein partitioning comprises partitioning the input file into the one or more input shards based on different types of data.

17. The non-transitory computer readable medium of claim 16 , wherein the different types of data is dependent on data size.

18. The method of claim 1 , wherein the third process distributes the one or more input shards to the first process.

19. The system of claim 7 , wherein the third process distributes the one or more input shards to the first process.

20. The non-transitory computer readable medium of claim 13 , wherein the third process distributes the one or more input shards to the first process.

Assignments (2)
CHANGE OF NAME Recorded Mar 7, 2022
From: GOOGLE INC.
To: GOOGLE LLC
Reel/Frame 059334/0572 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 4, 2022
From: KONERDING, DAVID; BRECKENRIDGE, JORDAN M.; BELOV, DANIEL
To: GOOGLE INC.
Reel/Frame 059173/0373 →
Continuity (4)
Continuation 16229886 · Dec 21, 2018
Continuation 15362177 · Nov 28, 2016
Continuation 13432117 · Mar 28, 2012
Provisional Application 61468417 · Mar 28, 2011
References Cited (78)
US 5349682A · Rosenberry · 1994 [cited by applicant]
US 5901312A · Radko · 1999 [cited by examiner]
US 5963965A · Vogel · 1999 [cited by applicant]
US 6108683A · Kamada et al. · 2000 [cited by applicant]
US 6178519B1 · Tucker · 2001 [cited by applicant]
US 6408342B1 · Moore et al. · 2002 [cited by applicant]
US 6779016B1 · Aziz et al. · 2004 [cited by applicant]
US 7805407B1 · Verbeke · 2010 [cited by examiner]
US 7900206B1 · Joshi et al. · 2011 [cited by applicant]
US 7921363B1 · Hao · 2011 [cited by examiner]
US 8280955B1 · Tyagi · 2012 [cited by applicant]
US 8381212B2 · Brelsford · 2013 [cited by examiner]
US 8392482B1 · McAlister et al. · 2013 [cited by applicant]
US 8424082B2 · Chen et al. · 2013 [cited by applicant]
US 20020065943A1 · Czajkowski et al. · 2002 [cited by applicant]
US 20020194184A1 · Baskins et al. · 2002 [cited by applicant]
US 20030051169A1 · Sprigg · 2003 [cited by examiner]
US 20040205759A1 · Oka · 2004 [cited by applicant]
US 20050097252A1 · Kelley · 2005 [cited by examiner]
US 20060048157A1 · Dawson et al. · 2006 [cited by applicant]
US 20060167966A1 · Kumar et al. · 2006 [cited by applicant]
US 20060230405A1 · Fraenkel et al. · 2006 [cited by applicant]
US 20070050393A1 · Vogel et al. · 2007 [cited by applicant]
US 20070156907A1 · Galchev · 2007 [cited by examiner]
US 20070180451A1 · Ryan et al. · 2007 [cited by applicant]
US 20070294319A1 · Mankad et al. · 2007 [cited by applicant]
US 20070294697A1 · Theimer et al. · 2007 [cited by applicant]
US 20080007765A1 · Ogata et al. · 2008 [cited by applicant]
US 20080115143A1 · Shimizu et al. · 2008 [cited by applicant]
US 20080155103A1 · Bailey · 2008 [cited by applicant]
US 20080177424A1 · Wheeler · 2008 [cited by examiner]
US 20080295112A1 · Muscarella · 2008 [cited by applicant]
US 20080306761A1 · George et al. · 2008 [cited by applicant]
US 20090055469A1 · Burckart et al. · 2009 [cited by applicant]
US 20090282477A1 · Chen · 2009 [cited by applicant]
US 20100017461A1 · Kokkevis · 2010 [cited by applicant]
US 20100030800A1 · Brodfuehrer et al. · 2010 [cited by applicant]
US 20100115097A1 · Gnanasambandam et al. · 2010 [cited by applicant]
US 20100125847A1 · Hayashi · 2010 [cited by applicant]
US 20100186018A1 · Bell, Jr. et al. · 2010 [cited by applicant]
US 20100269120A1 · Manttari et al. · 2010 [cited by applicant]
US 20110191781A1 · Karanam et al. · 2011 [cited by applicant]
US 20110302583A1 · Abadi et al. · 2011 [cited by applicant]
US 20120084798A1 · Reeves · 2012 [cited by applicant]
US 20120136835A1 · Kosuru et al. · 2012 [cited by applicant]
US 20120166514A1 · Mathew · 2012 [cited by examiner]
US 20130138768A1 · Huang · 2013 [cited by applicant]
US 20130139162A1 · Jackson · 2013 [cited by applicant]
US 20130151707A1 · Boutin et al. · 2013 [cited by applicant]
US 20130174174A1 · Jung et al. · 2013 [cited by applicant]
US 20130212165A1 · Vermeulen et al. · 2013 [cited by applicant]
US 20130325997A1 · Higgins et al. · 2013 [cited by applicant]
US 20130332612A1 · Cai et al. · 2013 [cited by applicant]
US 20140013301A1 · Duggal et al. · 2014 [cited by applicant]
US 20140025774A1 · Gosain et al. · 2014 [cited by applicant]
210 DAGMan Applications downloaded from the Internet on Mar. 28, 2011. http://www.cs.wisc.edu/condor/manual/v7.0/2 10 DAGMan Applications.html, 18 pages. [cited by applicant]
Amazon Elastic Compute Cloud (Amazon EC2) downloaded from the internet on Mar. 28, 2011. http://aws.amazon.com/ec2/ , 11 pages. [cited by applicant]
Amazon Elastic Compute Cloud API Reference API Version Feb. 28, 2011 by Amazon Web Services, downloaded from the internet on Mar. 28, 2011, http://awsdocs.s3.amazonaws.com/EC2/2011-02-28/ec2-api-2011-02-28.pdf, 414 page… [cited by applicant]
Amazon Elastic Compute Cloud Command Line Tools Reference API Version Feb. 28, 2011 by Amazon Web Services, downloaded from the internet on Mar. 28, 2011, http://awsdocs.s3.amazonaws.com/EC2/2011-02-28/ec2-clt-2011-02-2… [cited by applicant]
Amazon Elastic Compute Cloud Getting Started Guide API Version Feb. 28, 2011 by Amazon Web Services, downloaded from the internet on Mar. 28, 2011, http://awsdocs.s3.amazonaws.com/EC2/2011-02-28/ec2-gsg-2011-02-28.p-df,… [cited by applicant]
Amazon Elastic Compute Cloud User Guide API Version Feb. 28, 2011 by Amazon Web Services, downloaded from the internet on Mar. 28, 2011, http://awsdocs.s3.amazonaws.com/EC2/2011-02-28/ec2-ug-2011-02-28.pdf, 373 pages. [cited by applicant]
Cluster Resources Moab Workload Manager, downloaded from the internet on Mar. 28, 2011. http://www.clusterresources.com/products/moab-cluster-suite/workload-mana-ger.php, 2 pages. [cited by applicant]
Directed Acyclic Graph Manager, Condor Hihg Throughput Computing, DAGMan, downloaded from the internet on Mar. 28, 2011. http://www.cs.wisc.edu/condor/dagman/, 1 page. [cited by applicant]
GRAM—Globus downloaded from the internet on Mar. 28, 2011. http:/dev.globus.org/wiki/GRAM, 3 pages. [cited by applicant]
GRAM4 Approach Execution Management: Key Concepts, downloaded from the internet on Mar. 28, 2011, http://www.globus.org/toolkit/docs/4.2/4.2.1/execution/key/, 24 pages. [cited by applicant]
Madrid, Marcella, “How to Get Started Using Your TeraGrid Allocation”, PCS, Feb. 2011, downloaded from the internet on Mar. 28, 2011, http://www.teragridforum.org/mediawiki/images/1/15/NewUserTraining.sub.—TG10.pdf, 68 … [cited by applicant]
Moab Workload Manager Administrators Guide, Version 6.0.1, Adaptive Computing, downloaded on Mar. 28, 2011, http://www.adaptivecomputing.com/resources/docs/mwm/ , 841 pages. [cited by applicant]
Notice of Allowance issued in U.S. Appl. No. 13/432,088, filed Aug. 26, 2015, 18 pages. [cited by applicant]
Notice of Allowance issued in U.S. Appl. No. 13/432,117 dated Aug. 26, 2016, 14 pages. [cited by applicant]
Notice of Allowance issued in U.S. Appl. No. 15/362,177, mailed on Aug. 28, 2018, 15 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,070, filed Feb. 6, 2014, 27 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,074, filed Apr. 2, 2014, 16 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,074, filed Jul. 8, 2013, 18 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,088, filed Oct. 8, 2014, 30 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,117 dated Jul. 29, 2015, 22 pages. [cited by applicant]
Office Action issued in U.S. Appl. No. 13/432,117 dated Oct. 3, 2014, 28 pages. [cited by applicant]
TeraGrid [About] downloaded from the internet on Mar. 28, 2011. www.teragrid.org/web/about/index, 1 page. [cited by applicant]
Windows Azure downloaded from the internet on Mar. 28, 2011. http://www.microsoft.com/windowsazure/windowsazure/ 2 pages. [cited by applicant]