IP Library › Granted Patent US 9,218,220
Granted Patent B2
US 9,218,220 · App. 13/613,183 · Granted Dec 22, 2015

Elastic and scalable publish/subscribe service

Inventors: Han Chen (White Plains, NY); Minkyong Kim (Scarsdale, NY); Hui Lei (Scarsdale, NY); Ming Li (Elmsford, NY); Fan Ye (Yorktown Heights, NY)
Assignee: INTERNATIONAL BUSINESS MACHINES CORPORATION
G06F9/5083H04L67/26H04L67/1008H04L67/18
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 9,218,220
App. No.
13/613,183
Granted
Dec 22, 2015
Kind
B2
Abstract

A system and method are disclosed for an elastic and scalable publish/subscribe scheme. Subscription information is received at a dispatcher node. A plurality of matching nodes is selected in an overlay network to store the subscription information on a computer readable storage medium. Upon receiving an event at a dispatching node, at least one of the matching nodes with the stored subscription information is selected to process the event.

Claims (32)

1. A publish/subscribe system, comprising:

a plurality of matching nodes forming an overlay network, wherein the matching nodes are configured to match events to subscription information stored on a computer-readable storage medium at each of said matching nodes; and

at least one dispatcher node comprising:

a processor configured to forward subscription information and events received at the at least one dispatcher node to at least one of the matching nodes based on the matching node being a single hop from the at least one dispatcher node;

a load collector configured to periodically receive load information from the plurality of matching nodes; and

a partitioner configured to add and/or remove matching nodes from the overlay network based on the load information.

2. The system of claim 1 , wherein the at least one dispatcher node comprises a subscription forwarder configured to select a plurality of matching nodes for storing received subscription information.

3. The system of claim 2 , wherein the subscription forwarder employs a multi-dimensional subscription space partitioning technique which partitions each attribute associated with a subscription into a plurality of subscription space segments, and assigns each subscription space segment to a different matching node.

4. The system of claim 1 , wherein the at least one dispatcher node comprises an event forwarder configured to select at least one matching node to process an event received at the at least one dispatcher node.

5. The system of claim 4 , wherein the event forwarder selects the at least one matching node with a lowest predicted response time to execute a matching procedure to determine which events are to be transmitted to subscribers.

6. The system of claim 5 , wherein the lowest predicted response time is computed based on at least one of: length of an event queue, message rate or match rate.

7. The system of claim 1 , further comprising a gossiper configured to periodically transmit data associated with a matching node to a predetermined number of randomly selected matching nodes.

8. The system of claim 7 , wherein the data indicates a liveliness value, contact information and load information associated with a matching node.

9. A non-transitory computer readable storage medium comprising a computer readable program, wherein the computer readable program when executed on a computer causes the computer to perform the steps of:

receiving, by a dispatching node, subscription information at a dispatching node;

selecting, by the dispatching node, a plurality of matching nodes in an overlay network to store the subscription information on a computer readable storage medium at each of the plurality of matching nodes based on each of the plurality of matching nodes being a single hop from the dispatching node;

receiving an event at a dispatching node;

selecting at least one of the matching nodes with the stored subscription information to process the event;

periodically receiving, by the dispatching node, load information from each of the matching nodes; and

adding or removing, by the dispatching node, matching nodes from the overlay network, based on the load information.

10. A non-transitory computer readable storage medium comprising a computer readable program, wherein the computer readable program when executed on a computer causes the computer to perform the steps of:

receiving, by a dispatching node, subscription information at a dispatching node;

selecting, by the dispatching node, a plurality of matching nodes in an overlay network to store the subscription information on a computer readable storage medium at each of the plurality of matching nodes based on each of the plurality of matching nodes being a single hop from the dispatching node;

receiving an event at a dispatching node;

selecting the matching node with a lowest predicted response time to execute a matching procedure for determining which events are to be transmitted to subscribers;

periodically transmitting data associated with a matching node to a predetermined number of randomly selected matching nodes using a gossiping protocol;

periodically receiving, by the dispatching node, load information from each of the matching nodes; and

adding or removing, by the dispatching node, matching nodes from the overlay network, based on the load information.

11. The non-transitory computer readable storage medium of claim 10 , wherein selecting a plurality of matching nodes in an overlay network to store the subscription information includes applying a multi-dimensional subscription space partitioning technique to partition each attribute associated with the subscription information into a plurality of subscription space segments.

12. The non-transitory computer readable storage medium of claim 11 , further comprising assigning an entirety of the subscription to different matching nodes which are each responsible for at least one of the subscription space segments.

13. The non-transitory computer readable storage medium of claim 11 , further comprising identifying the subscription space segments that are associated with attribute values of an incoming event using information propagated through gossiping, and identifying a set of candidate matching nodes responsible for the identified subscriptions space segments.

14. The non-transitory computer readable storage medium of claim 10 , wherein the lowest predicted response time is computed based on at least one of: length of an event queue, message rate or match rate.

Continuity (2)
Continuation 13014501 · Jan 26, 2011
Related Publication 20130007131A1 · Jan 3, 2013