IP Library Granted Patent US 11,080,090
Granted Patent B2
US 11,080,090 · App. 15/503,889 · Granted Aug 3, 2021

Method and system for scalable job processing

Inventors: Matthew James George Painter (London, GB); Ian Andrew Clark (London, GB)
Assignee: IMPORT.IO LIMITED
G06F9/4881G06F9/485G06F9/5083G06F9/546
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 11,080,090
App. No.
15/503,889
Granted
Aug 3, 2021
Kind
B2
Abstract

The present invention relates to methods for processing jobs within a cluster architecture. One method comprises the pausing of a job when waiting upon external dependencies. Another method comprises the transmission of messages relating to the ongoing processing of jobs back to a client via a persistent messaging channel. Yet another method comprises determining capacity at a node before allocating a job for processing by the node or adding the job to a cluster queue. A system for processing jobs within a cluster architecture is also disclosed.

Claims (46)

1. A method for processing jobs in a cluster architecture including a cluster queue for providing jobs to a plurality of processing nodes, each node of the plurality of processing nodes including a local queue for providing jobs, the method including:

one node of a plurality of processing nodes within the cluster architecture receiving a job;

processing the job until the job is waiting for a dependency external to the cluster architecture to complete;

when the job is waiting for the external dependency, pausing the job by the node;

after pausing the job, retrieving another job from a local queue for the node if the local queue includes another job and from the cluster queue if the local queue does not include another job;

processing the retrieved another job by the node; and

when the external dependency is completed, allocating the paused job to the local queue for the node to continue the job;

wherein the job is received from an external client and messages in relation to the job are delivered back to the external client via a persistent messaging channel;

further comprising determining, by the node, whether to process the received job locally or to push the received job to the cluster queue based on a current and/or projected processing load of the node; and processing the job by the node based on determining to process the received job locally.

2. A method as claimed in claim 1 , wherein each job is deserialised from a job definition.

3. A method as claimed in claim 2 , wherein the job definition includes task parameters and definition type.

4. A method as claimed in claim 1 , wherein continued jobs are allocated to the local queue for the node.

5. A method as claimed in claim 1 , wherein, when a job is paused, its state is saved.

6. A method as claimed in claim 1 , wherein the external dependency is network input/output.

7. A method as claimed in claim 1 , wherein the external dependency is information requested from an external Internet service.

8. A method as claimed in claim 1 , wherein the node receives another job from another node of the plurality of processing nodes.

9. A system comprising:

a plurality of processing nodes within a cluster architecture including a cluster queue for providing jobs to the plurality of processing nodes, each node of the plurality of processing nodes including a local queue for providing jobs; and

a communications system;

wherein the system is configured to perform a method for processing jobs in the cluster architecture, the method comprising:

one node of a plurality of processing nodes within the cluster architecture receiving a job;

processing the job until the job is waiting for a dependency external to the cluster architecture to complete;

when the job is waiting for the external dependency, pausing the job by the node;

after pausing the job, retrieving another job from a local queue for the node if the local queue includes another job and from the cluster queue if the local queue does not include another job;

processing the received another job by the node; and

when the external dependency is completed, allocating the paused job to the local queue for the node to continue the job;

wherein the job is received from an external client and messages in relation to the job are delivered back to the external client via a persistent messaging channel;

further comprising determining, by the node, whether to process the received job locally or to push the received job to the cluster queue based on a current and/or projected processing load of the node; and processing the job by the node based on determining to process the received job locally.

10. A non-transitory computer readable storage medium having stored therein computer instructions which, when executed by a processor of a node of a plurality of processing nodes, cause the node to perform operations for processing jobs in a cluster architecture including a cluster queue for providing jobs to the plurality of processing nodes, each node of the plurality of processing nodes including a local queue for providing jobs, the operations comprising:

one node of the plurality of processing nodes within the cluster architecture receiving a job;

processing the job until the job is waiting for a dependency external to the cluster architecture to complete;

when the job is waiting for the external dependency, pausing the job by the node;

after pausing the job, retrieving another job from a local queue for the node if the local queue includes another job and from the cluster queue if the local queue does not include another job;

processing the retrieved another job by the node; and

when the external dependency is completed, allocating the paused job to the local queue for the node to continue the job;

wherein the job is received from an external client and messages in relation to the job are delivered back to the external client via a persistent messaging channel;

further comprising determining, by the node, whether to process the received job locally or to push the received job to the cluster queue based on a current and/or projected processing load of the node; and processing the job by the node based on determining to process the received job locally.

11. A method as claimed in claim 1 , wherein the external dependency is non-deterministic.

12. A method as claimed in claim 1 , wherein the messages include results and wherein at least some messages are delivered before the job is completed.

13. A method as claimed in claim 1 , wherein the local queue is stored within memory at the one node.

14. A method as claimed in claim 1 , wherein the local queue stores serialised versions of jobs.

15. The system of claim 9 , wherein the external dependency is non-deterministic.

16. The system of claim 9 , wherein the messages include results and wherein at least some messages are delivered before the job is completed.

17. The system of claim 9 , wherein the local queue is stored within memory at the one node.

18. The system of claim 9 , wherein the local queue stores serialised versions of jobs.

19. The non-transitory computer readable storage medium of claim 10 , wherein the external dependency is non-deterministic.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 18, 2018
From: PAINTER, MATTHEW JAMES GEORGE; CLARK, IAN ANDREW
To: IMPORT.IO LIMITED
Reel/Frame 045571/0047 →
Priority Claims (1)
GB 1414463 · Aug 14, 2014 · national
Continuity (1)
Related Publication 20180349178A1 · Dec 6, 2018