IP Library Granted Patent US 7,380,005
Granted Patent B1
US 7,380,005 · App. 11/943,150 · Granted May 27, 2008

Systems, methods and computer program products for improving placement performance of message transforms by exploiting aggressive replication

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 7,380,005
App. No.
11/943,150
Granted
May 27, 2008
Kind
B1
Abstract

Systems, methods and computer program products for improving overall end-to-end runtime latency of flow graphs of message transformations which are placed onto an overlay network of broker machines by aggressively replicating stateless transformations. Exemplary embodiments include a method including defining a message transformation graph, receiving information about measured and estimated properties of a message flow, receiving information about physical brokers and links in the overlay network onto which the message transformation graph is deployed, labeling each of a plurality of stateless transformations associated with the flow graph as replicable, heuristically determining a number of replicas and a corresponding load partitioning ratios among the number of replicas for each stateless transformation, converting the message transformation graph into an enhanced flow graph, running a placement algorithm with the enhanced flow graph and consolidating each of the plurality of virtual replicas that are assigned to a common message broker.

Claims (13)

1. In a computer system coupled to an overlay network of broker machines, a method for placing message transformations on the broker machines to improve end-to-end performance, the method consisting of:

defining a message transformation flow graph including computational nodes and edges;

receiving information about measured and estimated properties of a message flow associated with the message transformation flow graph;

receiving information about physical brokers and links in the overlay network onto which the message transformation flow graph is deployed, wherein the information about the physical broker and links are selected from a group consisting of: CPU capacity of each of the physical brokers, link capacity and link latency;

labeling each of a plurality of stateless transformations associated with the message transformation flow graph as replicable;

heuristically determining a number of replicas and a corresponding load partitioning ratios among the number of replicas for each of the plurality of replicable stateless transformations;

converting the message transformation flow graph into an enhanced flow graph having a plurality of virtual replicas of each of the plurality of replicable stateless transformations, and having a plurality of additional data partitioning filter transformations configured to partition a workload for each of the plurality of stateless transformations labeled as replicable;

running a placement algorithm with the enhanced flow graph to generate an optimal assignment of the message transformations in the enhanced flow graph to the broker machines in the overlay network; and

consolidating each of the plurality of virtual replicas that are assigned to a common message broker.

2. The method as claimed in claim 1 wherein the information about measured and estimated properties of a message flow associated with the flow graph, are selected from a group consisting of: message rates into a message transformation, message size, CPU utilization per message in each transformation, and expected number of output messages per input message.

3. The method as claimed in claim 2 wherein the message transformation flow graph is converted into the enhanced flow graph by replicating the plurality of replicable stateless transformations with the heuristically determined number of replicas and the corresponding load partitioning ratios among the number of replicas for each of the plurality of replicable stateless transformations.

4. The method as claimed in claim 3 wherein the placement algorithm uses an objective function based on expected overall end-to-end latency between producers and consumers.

5. The method as claimed in claim 4 wherein the generated optimal assignment assigns the each of the plurality of replicable stateless transformations to a particular broker, and replicas of each stateful transformation to one or more brokers, together with a percentage of message traffic destined for the each of the plurality of replicable stateless transformations that is to be routed to the replicas of the each of the plurality of replicable stateless transformations.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jun 17, 2016
From: INTERNATIONAL BUSINESS MACHINES CORPORATION
To: HULU, LLC
Reel/Frame 039071/0323 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 20, 2007
From: LI, YING; STROM, ROBERT
To: INTERNATIONAL BUSINESS MACHINES CORPORATION
Reel/Frame 020140/0486 →