IP Library › Granted Patent US 10,469,396
Granted Patent B2
US 10,469,396 · App. 14/879,679 · Granted Nov 5, 2019

Event processing with enhanced throughput

Inventors: David Mellor (Lynnfield, MA); Ora Lassila (Hollis, NH)
Assignee: PegaSystems, Inc.
H04L47/522G06Q10/00G06Q10/10G06Q10/107H04L45/02H04L49/90H04L67/10
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 10,469,396
App. No.
14/879,679
Granted
Nov 5, 2019
Kind
B2
Abstract

The present systems and methods allow for rapid processing of large volumes of events. A producer node in a cluster determines a sharding key for a received event from an event stream. The producer node uses a sharding map to correlate the sharding key for the event with a producer channel, and provides the event to a producer event buffer associated with the producer channel. The producer event buffer transmits the event to a corresponding consumer event buffer associated with a consumer channel on a consumer node. The event processing leverages a paired relationship between producer channels on the producer node and consumer channels on the consumer node, so as to generate enhanced throughput. The event processing also supports dynamic rebalancing of the system in response to adding or removing producer or consumer nodes, or adding or removing producer or consumer channels to or from producer or consumer nodes.

Claims (74)

1. A digital data processing system comprising

a producer node in communicative coupling with one or more consumer nodes and with a sharding map, wherein the producer node is configured to:

receive at least one event stream comprising a plurality of events;

determine a sharding key associated with an event among the plurality of events in the event stream;

identify, based on the sharding map, a producer event buffer associated with a producer channel on the producer node for transmitting the event to a corresponding consumer event buffer associated with a consumer channel on a consumer node among the one or more consumer nodes, wherein the sharding map correlates the sharding key for the event with the producer channel; and

provide the event to the producer event buffer associated with the producer channel in order to transmit the event to the corresponding consumer event buffer associated with the consumer channel on the consumer node,

wherein the producer node is further configured to initialize a plurality of producer channels by

referencing a channel map that correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes;

creating at least one producer channel based on the channel map that is communicatively coupled with a corresponding consumer channel among the one or more consumer channels; and

updating the sharding map to correlate the sharding key for the event with the at least one producer channel created based on the channel map.

2. The system of claim 1 ,

wherein the sharding map correlates the sharding key for the event with the producer channel based on a partition space, and

wherein the partition space is determined using a partition criterion based on a count of consumer nodes available to process the plurality of events in the event stream.

3. The system of claim 2 , wherein the producer node is further configured to update the sharding map in response to detecting an update to a channel map that correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes.

4. The system of claim 3 , wherein the system is configured to update the sharding map by

redistributing the partition space based on determining an updated partition size for existing producer channels based on a count of available consumer channels, and

assigning an overflow portion of the partition space based on the updated partition size to a new producer channel.

5. The system of claim 3 , wherein the producer node is configured to update the sharding map by redistributing the partition space based on

querying the one or more consumer nodes to identify resources available to the one or more consumer nodes, and

weighting the partitions assigned to each producer channel based on the identified resources.

6. The system of claim 1 , wherein the producer node is further configured to adjust the rate of providing events to the producer event buffer so as to manipulate the rate of events processed on the corresponding consumer node by the consumer event buffer, in order to improve throughput between the producer node and the consumer node.

7. The system of claim 1 , wherein the producer node is further configured to

bundle event transport metadata with the plurality of events in the event stream for transmission to the consumer node, wherein the event transport metadata contains instructions for identifying a rule to process the plurality of events after receipt by the consumer node, and

trigger the consumer node to execute the identified rule for processing the plurality of events.

8. The system of claim 1 , wherein the producer node is further configured to update a plurality of producer channels in response to detecting an update to at least one of the sharding map and a channel map by the consumer node,

wherein the update to the at least one of the sharding map and the channel map comprises at least one of adding a consumer node to the system, removing a consumer node from the system, adding one or more consumer channels to the consumer node, and removing one or more consumer channels from the consumer node, and

wherein the channel map correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes.

9. The system of claim 8 ,

wherein the producer node is further configured to transmit a message to the one or more consumer nodes in response to detecting the update to at least one of the sharding map and the channel map, and

wherein the transmitted message triggers the one or more consumer nodes to:

determine a delta that tracks the update to the at least one of the sharding map and the channel map,

identify data to be moved to at least one of a different consumer channel and a different consumer node, and

copy the data to be moved to a cluster change data map, wherein the data copied to the cluster change data map triggers the consumer node to

copy the moved data from the cluster change data map,

clear the copied data from the cluster change data map, and

update a status of the consumer node to active in a consumer map that tracks the one or more consumer nodes in the cluster, and

wherein the producer node is further configured to resume transmission of data to the consumer node in response to detecting the updated status of the consumer node in the consumer map.

10. A method for event processing, the method comprising:

receiving, by a producer node, at least one event stream comprising a plurality of events, wherein the producer node is in communicative coupling with one or more consumer nodes and with a sharding map;

determining, with the producer node, a sharding key associated with an event among the plurality of events in the event stream;

identifying, based on the sharding map, a producer event buffer associated with a producer channel on the producer node for transmitting the event to a corresponding consumer event buffer associated with a consumer channel on a consumer node among the one or more consumer nodes, wherein the sharding map correlates the sharding key for the event with the producer channel;

providing the event to the producer event buffer associated with the producer channel in order to transmit the event to the corresponding consumer event buffer associated with the consumer channel on the consumer node;

initializing a plurality of producer channels by:

referencing, with the producer node, a channel map that correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes;

creating at least one producer channel based on the channel map that is communicatively coupled with a corresponding consumer channel among the one or more consumer channels; and

updating the sharding map to correlate the sharding key for the event with the at least one producer channel created based on the channel map.

11. The method of claim 10 ,

wherein the sharding map correlates the sharding key for the event with the producer channel based on a partition space, and

wherein the partition space is determined using a partition criterion based on a count of consumer nodes available to process the plurality of events in the event stream.

12. The method of claim 11 , further comprising updating, with the producer node, the sharding map in response to detecting an update to a channel map that correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes.

13. The method of claim 12 , wherein the updating the sharding map comprises redistributing the partition space by:

determining, with the producer node, an updated partition size for existing producer channels based on a count of available consumer channels, and

assigning an overflow portion of the partition space based on the updated partition size to a new producer channel.

14. The method of claim 12 , wherein the updating the sharding map comprises redistributing the partition space by:

querying the one or more consumer nodes to identify resources available to the one or more consumer nodes, and

weighting the partitions assigned to each producer channel based on the identified resources.

15. The method of claim 10 , further comprising adjusting, with the producer node, the rate of providing events to the producer event buffer so as to manipulate the rate of events processed on the corresponding consumer node by the consumer event buffer, in order to improve throughput between the producer node and the consumer node.

16. The method of claim 10 , further comprising:

bundling, using the producer node, event transport metadata with the plurality of events in the event stream for transmission to the consumer node, wherein the event transport metadata contains instructions for identifying a rule to process the plurality of events after receipt by the consumer node, and

triggering the consumer node to execute the identified rule for processing the plurality of events.

17. The method of claim 10 , further comprising updating a plurality of producer channels in response to detecting an update to at least one of the sharding map and a channel map by the consumer node,

wherein the update to the at least one of the sharding map and the channel map comprises at least one of adding a consumer node to the system, removing a consumer node from the system, adding one or more consumer channels to the consumer node, and removing one or more consumer channels from the consumer node, and

wherein the channel map correlates an event in the event stream with one or more consumer channels on the one or more consumer nodes.

18. The method of claim 17 , further comprising:

transmitting, using the producer node, a message to the one or more consumer nodes in response to detecting an the update to at least one of the sharding map and the topic channel map,

wherein the transmitted message triggers the one or more consumer nodes to:

determine a delta that tracks the update to the at least one of the sharding map and the topic channel map,

identify data to be moved to at least one of a different consumer channel and a different consumer node, and

copy the data to be moved to a cluster change data map,

wherein the data copied to the cluster change data map triggers the consumer node to

copy the moved data from the cluster change data map,

clear the copied data from the cluster change data map, and

update a status of the consumer node to active in a consumer map that tracks the one or more consumer nodes in the cluster, and

resuming transmission, using the producer node, of data to the consumer node in response to detecting the updated status of the consumer node in the consumer map.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 11, 2016
From: MELLOR, DAVID; LASSILA, ORA
To: PEGASYSTEMS INC.
Reel/Frame 038241/0076 →
Continuity (2)
Provisional Application 62062515 · Oct 10, 2014
Related Publication 20160105370A1 · Apr 14, 2016
Cited By (1)
US 12,314,904