IP Library Granted Patent US 12,333,340
Granted Patent B1
US 12,333,340 · App. 17/515,117 · Granted Jun 17, 2025

Data processing pipeline horizontal scaling

Inventors: Prakul Agarwal (San Francisco, CA); Khwaja Mustafa Sidiqi (San Diego, CA); Stefanie Diem Anh Tonnu (San Jose, CA); Christian Yang (Sunnyvale, CA)
Assignee: Zoox, Inc.
G06F9/5027G06F9/4887G06F9/5088
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,333,340
App. No.
17/515,117
Granted
Jun 17, 2025
Kind
B1
Abstract

Techniques are disclosed for executing a data processing pipeline. The techniques may include receiving a job at a data pipeline queue, setting up one or more distributed processing environments, and allocating the job to one of the distributed processing environments. The techniques may further include receiving the allocated job at a job queue within the distributed processing environment, increasing a priority level of the job, and executing the job within the distributed processing environment. The techniques can further include providing a retry pipeline at the data processing pipeline, and re-executing the job at a stage following a failure of at least one of its components. The techniques may decrement the retry budget as the job is re-executed.

Claims (90)

1. A system comprising:

one or more processors; and

one or more non-transitory computer-readable media storing computer-executable instructions that, when executed, cause the one or more processors to perform operations comprising:

receiving, at an execution queue, a processing job associated with an operation of an autonomous vehicle;

determining that a scheduling component associated with a first distributed processing environment is associated with an overload attribute;

determining, by a controller, an overload state associated with the overload attribute;

generating, based at least in part on the overload state, a second distributed processing environment;

allocating, by the controller, the processing job to the second distributed processing environment;

determining, by the controller, that an actual completion time for the processing job is longer than an estimated completion time for the processing job; and

generating, by the controller, a third distributed processing environment.

2. The system of claim 1 , wherein the first distributed processing environment comprises:

the scheduling component; and

a plurality of computing resources allocated to the first distributed processing environment.

3. The system of claim 1 , wherein the overload attribute comprises at least one of:

a first number of processing jobs executing at the first distributed processing environment;

a second number of processing jobs within the execution queue;

a total processing job runtime;

an organizational attribution associated with the first distributed processing environment;

a processing job response time;

a processing job turnaround time; or

a wait time associated with the execution queue.

4. The system of claim 1 , wherein the first distributed processing environment is configured to operate independently from the second distributed processing environment.

5. The system of claim 1 , the operations further comprising:

allocating, by the controller, a second processing job to the first distributed processing environment or the second distributed processing environment based on at least one of:

a first load level of the first distributed processing environment and a second load level of the second distributed processing environment;

an execution constraint of the second processing job; or

a first hardware availability of the first distributed processing environment and a second hardware availability of the second distributed processing environment.

6. A method comprising:

receiving, at an execution queue, a processing job;

determining that a scheduling component of a first distributed processing environment is associated with an attribute;

determining, by a controller, an unavailable state associated with the attribute;

generating, based at least in part on the unavailable state, a second distributed processing environment;

allocating, by the controller, the processing job to the second distributed processing environment;

determining, by the controller, that an actual completion time for the processing job is longer than an estimated completion time for the processing job; and

generating, by the controller, a third distributed processing environment.

7. The method of claim 6 , wherein the first distributed processing environment comprises:

the scheduling component; and

a plurality of computing resources allocated to the first distributed processing environment.

8. The method of claim 6 , wherein the unavailable state is indicative of an overload state of the scheduling component.

9. The method of claim 6 , wherein the attribute comprises at least one of:

a first number of processing jobs executing at the first distributed processing environment;

a second number of processing jobs within the execution queue;

a total processing job runtime;

an organizational attribution associated with the first distributed processing environment;

a processing job response time;

a processing job turnaround time; or

a wait time associated with the execution queue.

10. The method of claim 6 , further comprising:

transmitting, based at least in part on the unavailable state of a component associated with the attribute at the first distributed processing environment, the component associated with the attribute from the second distributed processing environment to the first distributed processing environment.

11. The method of claim 6 , wherein the first distributed processing environment is configured to operate independently of the second distributed processing environment.

12. The method of claim 6 , further comprising:

receiving, at the execution queue, a second processing job;

determining that the second processing job is associated with a local processing component;

determining, by the controller, an unavailability state associated with the local processing component of the first distributed processing environment; and

receiving, at the execution queue and by the controller, a third processing job associated with an operation of a vehicle.

13. The method of claim 6 , further comprising:

allocating, by the controller, a second processing job to the first distributed processing environment or the second distributed processing environment based on at least one of:

a first load level of the first distributed processing environment and a second load level of the second distributed processing environment;

an execution constraint of the second processing job; or

a first hardware availability of the first distributed processing environment and a second hardware availability of the second distributed processing environment.

14. The method of claim 6 , further comprising:

receiving the estimated completion time for the processing job.

15. The method of claim 6 , further comprising:

receiving a threshold associated with a maximum amount of jobs executable at the second distributed processing environment; and

determining, by the controller, that an amount of jobs to be executed at the second distributed processing environment meets or exceeds the threshold.

16. One or more non-transitory computer-readable media storing instructions executable by a processor, wherein the instructions, when executed, cause the processor to perform operations comprising:

receiving, at an execution queue, a processing job;

determining that a scheduler component of a first distributed processing environment is associated with an attribute;

determining, by a controller, an unavailable state associated with the attribute;

generating, based at least in part on the unavailable state, a second distributed processing environment;

allocating, by the controller, the processing job to the second distributed processing environment;

determining, by the controller, that an amount of jobs to be executed at the second distributed processing environment meets or exceeds a threshold associated with a maximum amount of jobs executable at the second distributed processing environment; and

generating, by the controller, a third distributed processing environment.

17. The one or more non-transitory computer-readable media of claim 16 , wherein the first distributed processing environment comprises:

the scheduler component; and

a plurality of computing resources allocated to the first distributed processing environment.

18. The one or more non-transitory computer-readable media of claim 16 , wherein the unavailable state is indicative of an overload state of the scheduler component.

19. The one or more non-transitory computer-readable media of claim 16 , wherein the attribute comprises at least one of:

a first number of processing jobs executing at the first distributed processing environment;

a second number of processing jobs within the execution queue;

a total processing job runtime;

an organizational attribution associated with the first distributed processing environment;

a processing job response time;

a processing job turnaround time; or

a wait time associated with the execution queue.

20. The one or more non-transitory computer-readable media of claim 16 , the operations further comprising:

transmitting, by the controller, a second processing job to the first distributed processing environment or the second distributed processing environment based on at least one of:

a first load level of the first distributed processing environment and a second load level of the second distributed processing environment;

an execution constraint of the second processing job; or

a first hardware availability of the first distributed processing environment and a second hardware availability of the second distributed processing environment.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 29, 2021
From: AGARWAL, PRAKUL; SIDIQI, KHWAJA MUSTAFA; TONNU, STEFANIE DIEM ANH; YANG, CHRISTIAN
To: ZOOX, INC.
Reel/Frame 057968/0528 →
References Cited (24)
US 7673190B1 · Engelbrecht et al. · 2010 [cited by applicant]
US 8205208B2 · Mausolf et al. · 2012 [cited by applicant]
US 9164808B2 · Parker · 2015 [cited by examiner]
US 10108920B2 · Boyacigiller et al. · 2018 [cited by applicant]
US 11175950B1 · Yang et al. · 2021 [cited by applicant]
US 11356508B1 · Vergara et al. · 2022 [cited by applicant]
US 11520642B2 · Kawano et al. · 2022 [cited by applicant]
US 11907764B2 · Drapkin et al. · 2024 [cited by applicant]
US 20050022094A1 · Mudge et al. · 2005 [cited by applicant]
US 20070214381A1 · Goyal et al. · 2007 [cited by applicant]
US 20100058107A1 · Blaauw et al. · 2010 [cited by applicant]
US 20110107166A1 · Flautner et al. · 2011 [cited by applicant]
US 20110161972A1 · Dillenberger et al. · 2011 [cited by applicant]
US 20130151819A1 · Piry et al. · 2013 [cited by applicant]
US 20160098292A1 · Boutin et al. · 2016 [cited by applicant]
US 20180129573A1 · Iturbe et al. · 2018 [cited by applicant]
US 20180210460A1 · Rowley · 2018 [cited by examiner]
US 20190065241A1 · Wong et al. · 2019 [cited by applicant]
US 20200012521A1 · Wu et al. · 2020 [cited by applicant]
US 20200142753A1 · Harwood et al. · 2020 [cited by applicant]
US 20220066813A1 · Taher et al. · 2022 [cited by applicant]
Office Action for U.S. Appl. No. 17/515,078, mailed on Jun. 27, 2024, Hermann et al., “Data Processing Pipeline Priority Allocation”, 38 pages. [cited by applicant]
Office Action for U.S. Appl. No. 17/515,078, dated Apr. 2, 2025, 29 pages. [cited by applicant]
Office Action for U.S. Appl. No. 17/515,078, dated Nov. 19, 2024, 28 pages. [cited by applicant]