IP Library Granted Patent US 7,689,863
Granted Patent B1
US 7,689,863 · App. 10/983,881 · Granted Mar 30, 2010

Restartable database loads using parallel data streams

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 7,689,863
App. No.
10/983,881
Granted
Mar 30, 2010
Kind
B1
Abstract

A method and computer program for reducing the restart time for a parallel application are disclosed. The parallel application includes a plurality of parallel operators. The method includes repeating the following: setting a time interval to a next checkpoint; waiting until the time interval expires; sending checkpoint requests to each of the plurality of parallel operators; and receiving and processing messages from one or more of the plurality of parallel operators. The method also includes receiving a checkpoint. request message on a control data stream, waiting to enter a state suitable for checkpointing, and sending a response message on the control data stream.

Claims (82)

1. A method for reducing the restart time for a parallel application, the parallel application including a plurality of parallel operators, the method comprising:

repeating the following by a CRFC component separate from the parallel operators:

setting a time interval to a next checkpoint;

waiting until the time interval expires;

sending checkpoint requests to each of the plurality of parallel operators;

receiving and processing messages from one or more of the plurality of parallel operators; and

performing the following by the parallel operators:

receiving a checkpoint request message;

waiting to enter a state suitable for checkpointing before taking a checkpoint; and

sending a response message.

2. The method of claim 1 further comprising: before entering the repeat loop: receiving a ready message by the CRFC component from each of the plurality of parallel operators indicating the parallel operator that originated the message is ready to accept checkpoint requests.

3. The method of claim 1 wherein receiving and processing messages from one or more of the plurality of parallel operators comprises:

receiving a checkpoint information message, including checkpoint information, from one of the plurality of parallel operators; and

storing the checkpoint information, along with an identifier for the one of the parallel operators, in a checkpoint data store.

4. The method of claim 1 wherein receiving and processing messages from one or more of the plurality of parallel operators comprises:

receiving a ready to proceed message from one of the plurality of parallel operators;

marking the one of the plurality of parallel operators as ready to proceed; and

if all of the plurality of parallel operators has been marked as ready to proceed, marking a current checkpoint as good.

5. The method of claim 1 where receiving and processing messages from one or more of the plurality of parallel operators comprises:

receiving a non-recoverable error message from one of the plurality of parallel operators; and

sending terminate messages to the plurality of parallel operators.

6. The method of claim 5 further comprising restarting the plurality of parallel operators.

7. The method of claim 6 , where restarting comprises:

sending initiate restart messages to the plurality of parallel processors; and

processing restart messages from the plurality of parallel processors.

8. The method of claim 7 where processing restart messages comprises:

receiving an information request message from one or more of the plurality of parallel operators;

retrieving checkpoint information regarding the one or more of the plurality of parallel operators from the checkpoint data store; and

sending the retrieved information to the one of the plurality of parallel operators.

9. The method of claim 7 where processing restart messages comprises:

receiving a ready to proceed message from one of the plurality of parallel operators;

marking the one of the plurality of parallel operators as ready to proceed; and

sending proceed messages to all of the plurality of parallel operators if all of the plurality of parallel operators have been marked as ready to proceed.

10. The method of claim 7 where processing restart messages comprises:

receiving an error message from one or more of the plurality of parallel operators; and

terminating the processing of the plurality of parallel operators.

11. A method for one of a plurality of parallel operators to record its state, the method comprising the one of the plurality of parallel operators:

receiving a checkpoint request message on a control data stream;

waiting to enter a state suitable for checkpointing before taking a checkpoint; and

sending a response message on the control data stream.

12. The method of claim 11 wherein waiting to enter a state suitable for checkpointing comprises:

receiving a checkpoint marker on an input data stream;

finishing writing data to an output data stream; and

sending a checkpoint marker on the output data stream.

13. The method of claim 11 wherein waiting to enter a state suitable for checkpointing comprises:

waiting for all of the parallel operator's outstanding input/output requests to be processed.

14. A computer program, stored on a tangible storage medium, for use in reducing the restart time for a parallel application, the parallel application comprising a plurality of parallel operators, the computer program comprising:

a CRCF component separate from the parallel operators which includes executable instructions that cause a computer to repeat the following:

set a time interval to a next checkpoint;

wait until the time interval expires;

send checkpoint requests to the plurality of parallel operators;

receive and process messages from one or more of the plurality of parallel operators; and

a plurality of parallel components, each of which is associated with one of the plurality of parallel operators, and each of which includes executable instructions that cause a computer to:

receive a checkpoint request message from the CRCF;

wait to enter a state suitable for checkpointing; and

send a checkpoint response message to the CRCF.

15. The computer program of claim 14 wherein each of the parallel components include executable instructions that cause a computer to:

determine that one of the parallel operators has experienced a non-recoverable error; and,

in sending a response message to the CRCF, the parallel component associated with the one parallel operator causes the computer to:

send a non-recoverable error message to the CRCF;

in receiving and processing messages from one or more of the plurality of parallel operators, the CRCF causes the computer to:

receive the non-recoverable error message; and

send stop processing messages to the plurality of parallel operators in response to the non-recoverable error message.

16. The computer program of claim 15 wherein the CRCF further includes executable instructions that cause the computer to:

send an initiate restart message to one of the plurality of parallel operators.

17. The computer program of claim 16 wherein

in response to the restart message from the CRCF, the parallel component associated with the one parallel operator causes the computer to:

send an information request to the CRCF;

in responding to the information request, the CRCF causes the computer to:

retrieve checkpoint information regarding the one parallel operator from a checkpoint data store; and

send the checkpoint information to the one parallel operator.

18. The computer program of claim 16 wherein

the parallel component associated with one of the parallel operators further comprises executable instructions that cause the computer to:

send a ready to proceed message to the CRCF;

in responding to the ready to proceed message, the CRCF causes the computer to:

mark the one parallel operator as ready to proceed; and

if all of the plurality of parallel operators have been marked as ready to proceed, sending proceed messages to all of the plurality of parallel operators.

19. The computer program of claim 16 wherein

the parallel component associated with one of the parallel operators further comprises executable instructions that cause the computer to:

send an error message to the CRCF;

in responding to the error message, the CRCF causes the computer to:

send messages to all of the parallel operators to terminate their processing.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 18, 2008
From: NCR CORPORATION
To: TERADATA US, INC.
Reel/Frame 020666/0438 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 8, 2004
From: KORENEVSKY, GREGORY; YUNG, ALEX P.
To: NCR CORPORATION
Reel/Frame 015978/0524 →