CA3078469C

Managing a computing cluster using time interval counters

Abstract

A method for processing state update requests in a distributed data processing system with a number of processing nodes includes maintaining a number of counters including a working counter indicating a current time interval, a replication counter indicating a time interval for which all requests associated with that time interval are replicated at multiple processing nodes of the number of processing nodes, and a persistence counter indicating a time interval of the number of time intervals for which all requests associated with that time interval are stored in persistent storage. The counters are used to manage processing of the state update requests.

CA3078469C, drawing sheet 1
Sheet 1 of 38

Term

12.1 yearsleft in the term

Expires 30 October 2038.

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

34 claims: 5 independent, 29 dependent

  1. 1
    A method for processing state update requests in a distributed data processing system including a plurality of processing nodes, the method including:processing a plurality of sets of requests using two or more of the plurality of processing nodes, each request of each set of requests being configured to cause a state update at a processing node of the plurality of processing nodes and being associated with a corresponding time interval of a plurality of time intervals, the plurality of sets of state update requests including a first set of requests associated with a first time interval of the plurality of time intervals;maintaining at a first processing node of the plurality of processing nodes, a plurality of counters, the plurality of counters including: a working counter indicating a current time interval, and a value thereof, of the plurality of time intervals in the distributed data processing system, a replication counter indicating a time interval, and a value thereof, of the plurality of time intervals for which all requests associated with that time interval are replicated at multiple processing nodes of the plurality of processing nodes, and a persistence counter indicating a time interval of the plurality of time intervals for which all requests associated with that time interval are stored in persistent storage associated with at least one processing node of the plurality of processing nodes, providing, a first message from the first processing node to the other processing nodes of the plurality of processing nodes at a first time, the first message including a first value of the working counter at the first time, a first value of the replication counter at the first time, and a first value of the persistence counter at the first time, wherein the first value of the replication counter lags the first value of the working counter and the first value of the persistence counter lags the first value of the .41. Date Reçue/Date Received 2021-09-20 replication counter, and wherein the replication counter in the first message indicates that all requests of a second set of requests of the plurality of sets of state update requests associated with an earlier time interval, prior to the first time interval, are replicated at two or more of the processing nodes in response to determining that the first replication value is equal or greater than the earlier time interval associated with the second set of requests, and any non-persistently stored requests associated with time intervals of the plurality of time intervals prior to the earlier time interval are replicated at two or more of the processing nodes, wherein at least some requests of the first set of requests associated with the first time interval are not yet replicated at two or more processing nodes of the plurality of processing nodes, and providing, a second message from the first processing node to the other processing nodes of the plurality of processing nodes at a second time subsequent to the first time, the second message includinga second value of the working counter at the second time, a second value of the replication counter at the second time, and a second value of the persistence counter at the second time, wherein the second value of the replication counter lags the second value of the working counter and the second value of the persistence counter lags the second value of the replication counter, and wherein the value of the replication counter in the second message indicates that all requests of the first set of requests associated with the first time interval are replicated at two or more of the processing nodes in response to determining that the second replication value is equal to or greater than the first time interval associated with the first set of requests, and any non-persistently stored requests associated with time intervals prior to the first time interval are replicated at two or more of the processing nodes, -42Date Reçue/Date Received 2021-09-20 wherein the second message causes at least one of the first set of processing nodes to complete storing operations for one or more requests of the first set of requests into persistent storage in response to determining that the second persistence value is equal to or greater than the first time interval associated with the first set of requests.
  2. 13
    An apparatus for processing state update requests, in a distributed data processing system including a plurality of processing nodes, the apparatus including:the distributed data processing system including the plurality of processing nodes;.44. Date Reçue/Date Received 2021-09-20 one or more processors included at two or more of the plurality of processing nodes for processing a plurality of sets of requests, each request of each set of requests being configured to cause a state update at a processing node of the plurality of processing nodes and being associated with a corresponding time interval of a plurality of time intervals, the plurality of sets of requests including a first set of requests associated with a first time interval of the plurality of time intervals;one or more data stores for maintaining at a first processing node of the plurality of processing nodes, a plurality of counters, the plurality of counters including: a working counter indicating a current time interval, and a value thereof, of the plurality of time intervals in the distributed data processing system, a replication counter indicating a time interval of the plurality of time intervals, and a value thereof, for which all requests associated with that time interval are replicated at multiple processing nodes of the plurality of processing nodes, and a persistence counter indicating a time interval, and a value thereof, of the plurality of time intervals for which all requests associated with that time interval are stored in persistent storage associated with at least one processing node of the plurality of processing nodes, a first output for providing a first message from the first processing node to the other processing nodes of the plurality of processing nodes at a first time, the first message including a first value of the working counter at the first time, a first value of the replication counter at the first time, and a first value of the persistence counter at the first time, wherein the first value of the replication counter lags the first value of the working counter and the first value of the persistence counter lags the first value of the replication counter, and wherein the replication counter in the first message indicates that -45Date Reçue/Date Received 2021-09-20 all requests of a second set of requests associated with an earlier time interval, prior to the first time interval, are replicated at two or more of the processing nodes in response to determining that the first replication value is equal or greater than the earlier time interval associated with the second set of requests, and any non-persistently stored requests associated with time intervals of the plurality of time intervals prior to the earlier time interval are replicated at two or more of the processing nodes, wherein at least some requests of the first set of requests associated with the first time interval are not yet replicated at two or more processing nodes of the plurality of processing nodes, and a second output for providing a second message from the first processing node to the other processing nodes of the plurality of processing nodes at a second time subsequent to the first time, the second message including a second value of the working counter at the second time, a second value of the replication counter at the second time, and a second value of the persistence counter at the second time, wherein the second value of the replication counter lags the second value of the working counter and the second value of the persistence counter lags the second value of the replication counter, and wherein the value of the replication counter in the second message indicates that all requests of the first set of requests associated with the first time interval are replicated at two or more of the processing nodes in response to determining that the second replication value is equal to or greater than the first time interval associated with the first set of requests, and any non-persistently stored requests associated with time intervals prior to the first time interval are replicated at two or more of the processing nodes, -46Date Reçue/Date Received 2021-09-20 wherein the second message causes at least one of the first set of processing nodes to complete storing operations for one or more requests of the first set of requests into persistent storage in response to determining that the second persistence value is equal to or greater than the first time interval associated with the first set of requests.
  3. 14
    A computing system for processing state update requests in a distributed data processing system including a plurality of processing nodes, the computing system including:means for processing a plurality of sets of requests using two or more of the plurality of processing nodes, each request of each set of requests being configured to cause a state update at a processing node of the plurality of processing nodes and being associated with a corresponding time interval of a plurality of time intervals, the plurality of sets of requests including a first set of requests associated with a first time interval of the plurality of time intervals;means for maintaining at a first processing node of the plurality of processing nodes, a plurality of counters, the plurality of counters including: a working counter indicating a current time interval, and a value thereof, of the plurality of time intervals in the distributed data processing system, a replication counter indicating a time interval, and a value thereof, of the plurality of time intervals for which all requests associated with that time interval are replicated at multiple processing nodes of the plurality of processing nodes, and a persistence counter indicating a time interval, and a value thereof, of the plurality of time intervals for which all requests associated with that time interval are stored in persistent storage associated with at least one processing node of the plurality of processing nodes, .47. Date Reçue/Date Received 2021-09-20 means for providing, a first message from the first processing node to the other processing nodes of the plurality of processing nodes at a first time, the first message including a first value of the working counter at the first time, a first value of the replication counter at the first time, and a first value of the persistence counter at the first time, wherein the first value of the replication counter lags the first value of the working counter and the first value of the persistence counter lags the first value of the replication counter, and wherein the replication counter in the first message indicates that all requests of a second set of requests associated with an earlier time interval, prior to the first time interval, are replicated at two or more of the processing nodes in response to determining that the first replication value is equal or greater than the earlier time interval associated with the second set of requests, and any non-persistently stored requests associated with time intervals of the plurality of time intervals prior to the earlier time interval are replicated at two or more of the processing nodes, wherein at least some requests of the first set of requests associated with the first time interval are not yet replicated at two or more processing nodes of the plurality of processing nodes, and means for providing, a second message from the first processing node to the other processing nodes of the plurality of processing nodes at a second time subsequent to the first time, the second message including a second value of the working counter at the second time, a second value of the replication counter at the second time, and a second value of the persistence counter at the second time, wherein the second value of the replication counter lags the second value of the working counter and the second value of the persistence counter lags the second value of the replication counter, and wherein the value of the replication counter in the second message indicates that -48Date Reçue/Date Received 2021-09-20 all requests of the first set of requests associated with the first time interval are replicated at two or more of the processing nodes in response to determining that the second replication value is equal to or greater than the first time interval associated with the first set of requests, and any non-persistently stored requests associated with time intervals prior to the first time interval are replicated at two or more of the processing nodes, wherein the second message causes at least one of the first set of processing nodes to complete storing operations for one or more requests of the first set of requests into persistent storage in response to determining that the second persistence value is equal to or greater than the first time interval associated with the first set of requests.
  4. 19
    The computing system according claim 17, wherein the means for aggregating the received first counts of state update requests for the first time interval and second counts of state update requests for the first time interval includes means for differencing a sum of the received first counts of state update requests for the first time interval and a sum of the second counts of state update requests for the first time interval.
  5. 27
    The apparatus according to any one of claims 13,25 and 26, further configured for receiving, at the first processing node, the first count of state update requests for the first time interval and the second count of state update requests for the first time interval from each processing node of the other processing nodes, aggregating the received first counts of state update requests for the first time interval and second counts of state update requests for the first time interval, and determining whether to increment the value of the replication counter based on the aggregation.