CA2392300C

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 check-point their state. Checkpointing in accordance with the invention make s 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.

CA2392300C, drawing sheet 1
Sheet 1 of 14

Term

Term ended

Expired 5 December 2020, 5.8 years ago.

  1. Priority
  2. Filed
  3. Granted
  4. Expired
  5. Today

47 claims: 6 independent, 41 dependent

  1. 1
    CA 02392300 2005-04-29 .76307-79 CLAIMS :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;and (b) checkpointing each process within the process stage in response to receipt by each process of at least one command message.
  2. 5
    A method for continuous flow checkpointing in a data processing system having at least one source for receiving and storing input data, at least one process for receiving and processing data from at least one source or prior process, and at least one sink for receiving processed data from the at least one process or source and for publishing processed data, the method including:(a) transmitting a checkpoint request message to every source;-20CA 02392300 2005-04-29 76307-79 (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 a 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 a 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;and (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 a state of such sink, saving any unpublished data, and resuming normal data processing in such sink.
  3. 24
    A method for continuous flow checkpointing in a data processing system having at least one source for receiving and storing input data, at least one process for receiving and processing data from at least one source or prior process, and at least one sink for receiving processed data from the at least one process or source and for publishing processed data, the method including:-24CA 02392300 2005-04-29 76307-79 (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 a 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 a 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 a 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. 25CA 02392300 2005-04-29 .76307-79
  4. 25
    Ά computer readable medium having computer readable instructions stored thereon, 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 readable instructions for causing a computer to:(a) propagate at least one command message through the process stage as part of the data flow;and (b) checkpoint each process within the process stage in response to receipt by each process of at least one command message.
  5. 29
    A computer readable medium having computer readable instructions stored thereon, for continuous flow checkpointing in a data, processing system having at least one source for receiving and storing input data, at least one process for receiving and processing data from at least one source or prior process, and at least one sink for -26CA 02392300 2005-04-29 .76307-79 receiving processed data from the at least one process or source and for publishing processed data, the computer readable 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 a 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 a 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;and (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 a state of such sink, save any unpublished data, and resume normal data processing in such sink.
  6. 47
    48. A computer readable medium having computer readable instructions stored thereon, for continuous flow checkpointing in a data processing system having at least one source for receiving and storing input data, at least one process for receiving and processing data from at least one source or prior process, and at least one sink for receiving processed data from at least one process or source and for publishing processed data, the computer readable 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;31CA 02392300 2005-04-29 .76307-79 (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. SMART & BIGGAR OTTAWA, CANADA PATENT AGENTS