CA2392300A1

Continuous flow checkpointing data processing

Abstract

A data processing system and method that provides checkpointing and permits a continuous flow of data processing by allowing each process to return to operation after checkpointing, independently of the time required by other processes to checkpoint their state. Checkpointing in accordance with the invention makes use of a command message from a checkpoint processor (300) that sequentially propagates through a process stage from data sources (302) through processes (304, 305) to data sinks (308), triggering each process to checkpoint its state and then pass on a checkpointing message to connected "downstream" processes. This approach provides checkpointing and permits a continuous flow of data processing by allowing each process to return to normal operation after checkpointing, independently of the time required by other processes to checkpoint their state.

CA2392300A1, drawing sheet 1
Sheet 1 of 12

Term

Term ended

Projected expiry passed 5 December 2020, 5.8 years ago.

  1. Priority
  2. Filed
  3. Published
  4. Projected expiry
  5. Today

47 claims: 6 independent, 41 dependent

  1. 1
    CA 02392300 2002-05-22 WO 01/42920 PCT7US00/42609 WHAT IS CLAIMED IS:1. A method for continuous flow checkpointing in a data processing system having at least one process stage comprising a data flow and at least two processes linked by the data flow, the method including: (a) propagating at least one command message through the process stage as part of the data flow;(b) checkpointing each process within the process stage in response to receipt by each process of at least one command message.
  2. 4
    5. A method for continuous flow checkpointing in a data processing system having one or more sources for receiving and storing input data, one or more processes for receiving and processing data from one or more sources or prior processes, and one or more sinks for receiving processed data from one or more processes or sources and for publishing processed data, the method including:(a) transmitting a checkpoint request message to every source;(b) suspending normal data processing in each source in response to receipt of such checkpoint request message, saving a current checkpoint record sufficient to reconstruct the state of such source, propagating a checkpoint message from such source to any process that consumes data from such source, and resuming normal data processing in each source;(c) suspending normal data processing in each process in response to receiving checkpoint messages from every source or prior process from which such process consumes data, saving a current checkpoint record sufficient to reconstruct the state of such process, propagating the checkpoint message from such process to any process or sink that consumes data from such process, and resuming normal data processing in such process;(d) suspending normal data processing in each sink in response to receiving checkpoint messages from every process from which such sink consumes data, saving a current checkpoint record sufficient to reconstruct the state of such sink, saving any unpublished data, and resuming normal data processing in such sink. -21CA 02392300 2002-05-22 WO 01/42920 PCT/US00/42609
  3. 23
    24. A method for continuous flow checkpointing in a data processing system having one or more sources for receiving and storing input data, one or more processes for receiving and processing data from one or more sources or prior processes, and one or more sinks for receiving processed data from one or more processes or sources and for publishing processed data, the method including:(a) transmitting a checkpoint request message to every source;(b) suspending normal data processing in each source in response to receipt of such checkpoint request message, saving a current checkpoint record sufficient to reconstruct the state of such source, propagating a checkpoint message from such source to any process that consumes data from such source, and resuming normal data processing in each source;(c) suspending normal data processing in each process in response to receiving checkpoint messages from every source or prior process from which such process consumes data, saving a current checkpoint record sufficient to reconstruct the state of such process, propagating the checkpoint message from such process to any process or sink that consumes data from such process, and resuming normal data processing in such process;(d) suspending normal data processing in each sink in response to receiving checkpoint messages from every process from which such sink consumes data, saving a current checkpoint record sufficient to reconstruct the state of such sink, saving any unpublished data, and propagating the checkpoint message from each sink to a checkpoint processor;(e) receiving the checkpoint messages from all sinks, and in response to such receipt, updating a stored value indicating completion of checkpointing in all sources, processes, and sinks, and transmitting the stored value to each sink;and (f) receiving the stored value in each sink and, in response to such receipt, publishing any unpublished data associated with such sink and resuming normal data processing in such sink. —25— CA 02392300 2002-05-22 WO 01/42920 PCT/US00/42609
  4. 24
    25. A computer program, stored on a computer-readable medium, for continuous flow checkpointing in a data processing system having at least one process stage comprising a data flow and at least two processes linked by the data flow, the computer program comprising instructions for causing a computer to:(a) propagate at least one command message through the process stage as part of the data flow;(b) checkpoint each process within the process stage in response to receipt by each process of at least one command message.
  5. 28
    29. A computer program, stored on a computer-readable medium, for continuous flow checkpointing in a data processing system having one or more sources for receiving and storing input data, one or more processes for receiving and processing data from one or more sources or prior processes, and one or more sinks for receiving processed data from one or more processes or sources and for publishing processed data, the computer program comprising instructions for causing a computer to:(a) transmit a checkpoint request message to every source;(b) suspend normal data processing in each source in response to receipt of such checkpoint request message, save a current checkpoint record sufficient to reconstruct the state of such source, propagate a checkpoint message from such source to any process that consumes data from such source, and resume normal data processing in each source;(c) suspend normal data processing in each process in response to receiving checkpoint messages from every source or prior process from which such process consumes data, save a current checkpoint record sufficient to reconstruct the state of such process, propagate the checkpoint message from such process to any process or sink that consumes data from such process, and resume normal data processing in such process;(d) suspend normal data processing in each sink in response to receiving checkpoint messages from every process from which such sink consumes data, save a current checkpoint record sufficient to reconstruct the state of such sink, save any unpublished data, and resume normal data processing in such sink. -27CA 02392300 2002-05-22 WO 01/42920 PCT/USOO/42609
  6. 47
    48. A computer program, stored on a computer-readable medium, for continuous flow checkpointing in a data processing system having one or more sources for receiving and storing input data, one or more processes for receiving and processing data from one or more sources or prior processes, and one or more sinks for receiving processed data from one or more processes or sources and for publishing processed data, the computer program comprising instructions for causing a computer to:(a) transmit a checkpoint request message to every source;(b) suspend normal data processing in each source in response to receipt of such checkpoint request message, save a current checkpoint record sufficient to reconstruct the state of such source, propagate a checkpoint message from such source to any process that consumes data from such source, and resume normal data processing in each source;(c) suspend normal data processing in each process in response to receiving checkpoint messages from every source or prior process from which such process consumes data, save a current checkpoint record sufficient to reconstruct the state of such process, propagate the checkpoint message from such process to any process or sink that consumes data from such process, and resume normal data processing in such process;(d) suspend normal data processing in each sink in response to receiving checkpoint messages from every process from which such sink consumes data, save a current checkpoint record sufficient to reconstruct the state of such sink, save any unpublished data, and propagate the checkpoint message from each sink to a checkpoint processor;(e) receive the checkpoint messages from all sinks, and in response to such receipt, update a stored value indicating completion of checkpointing in all sources, processes, and sinks, and transmit the stored value to each sink;and (f) receive the stored value in each sink and, in response to such receipt, publish any unpublished data associated with such sink and resume normal data processing in such sink.