IP Library Granted Patent US 7,673,065
Granted Patent B2
US 7,673,065 · App. 11/977,440 · Granted Mar 2, 2010

Support for sharing computation between aggregations in a data stream management system

Assignee: Oracle International Corporation
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,673,065
App. No.
11/977,440
Granted
Mar 2, 2010
Kind
B2
Abstract

A computer is programmed to process a continuous query that is known to perform a new aggregation on one or more stream(s) of data, using one or more other aggregations on the stream(s). The computer creates an operator to execute the continuous query, and schedules the operator for execution in a specific order. In several embodiments, the computer determines the order based on dependency of the new aggregation on other aggregation(s), and on the order of performance of the other aggregation(s). The new aggregation is scheduled for performance after performance of each of the other aggregations. The computer is further programmed to pass results of the other aggregations to the new aggregation, by execution of a predetermined function. Support for use of the other aggregations results within the new aggregation eliminates redundant computation of the other aggregations within the new aggregation. The new aggregation may be user defined or built-in.

Claims (48)

1. A method implemented in a computer of processing a plurality of streams of data, the method comprising:

processing the plurality of streams, to execute thereon a plurality of continuous queries based on a global plan;

during the processing, receiving a command to create a first aggregation and a first identification of a set of instructions comprising a function to be executed to perform the first aggregation;

during the processing, receiving a second identification of a second aggregation to be used by the first aggregation;

during the processing, creating in a memory of the computer, a first structure comprising the first identification and the second identification;

during the processing, receiving a new continuous query to be executed using the first aggregation;

during the processing, based on the first structure, creating in the memory an operator comprising at least one second structure, the second structure comprising a first field to hold a reference to the instance, and at least one additional field corresponding to at least one argument of the first aggregation;

during the processing, building in the memory a graph representing a plurality of aggregations as nodes;

wherein the plurality of aggregations comprises said first aggregation and said second aggregation;

wherein a plurality of directed edges are included in the graph, between each node and a group of nodes whose results are used by said each node;

during the processing, performing a sort on the graph, and using a result of the sort to store, in the memory, a temporal order for executing the plurality of aggregations;

during the processing, modifying the global plan in the memory by adding thereto the operator based on the temporal order, thereby to obtain a modified plan in the memory; and

altering the processing, to cause execution of the new continuous query in addition to the plurality of continuous queries, based on the modified plan in the memory;

during execution of the new continuous query creating in the memory an instance of the set of instructions;

invoking the function in the instance to process a tuple of the data received in a message, wherein the function is identified based at least on a type of the message; and

releasing in the memory, space occupied by the instance, in response to a predetermined condition being met.

2. The method of claim 1 wherein:

the second identification is received in said command, adjacent to the first Identification.

3. The method of claim 1 further comprising:

receiving with the command, identification of a value of a data type of the at least one argument;

wherein the function is further identified, for use in the invoking, based at least on the value of the data type.

4. The method of claim 1 wherein:

the first structure further comprises a first name of the first aggregation and a second name of the second aggregation; and

the first structure further comprises at least another data type of value to be returned by the aggregation.

5. A computer-readable storage medium comprising instructions to a computer to process a plurality of streams of data, the instructions comprising:

instructions to process the plurality of streams, to execute thereon a plurality of continuous queries based on a global plan;

instructions to receive a command to create a first aggregation and a first identification of a set of additional instructions comprising a function to be executed to perform the first aggregation;

instructions to receive a second identification of a second aggregation to be used by the first aggregation;

instructions to create in a memory of the computer, a first structure comprising the first identification and the second identification;

instructions to receive a new continuous query to be executed using the first aggregation;

instructions, based on the first structure, to create in the memory an operator comprising at least one second structure, the second structure comprising a first field to hold a reference to the instance, and at least one additional field corresponding to at least one argument of the first aggregation;

instructions to build in the memory a graph representing a plurality of aggregations as nodes;

wherein the plurality of aggregations comprises said first aggregation and said second aggregation;

wherein a plurality of directed edges are included in the graph, between each node and a group of nodes whose results are used by said each node;

instructions to perform a sort on the graph, and use a result of the sort to store, in the memory, a temporal order for executing the plurality of aggregations;

instructions to modify the global plan in the memory by adding thereto the operator based on the temporal order, thereby to obtain a modified plan in the memory; and

instructions to cause execution of the new continuous query in addition to the plurality of continuous queries, based on the modified plan in the memory;

instructions to create in the memory an instance of the set of additional instructions;

instructions to invoke the function in the instance to process a tuple of the data received in a message, wherein the function is identified based at least on a type of the message; and

instructions to release in the memory, space occupied by the instance, in response to a predetermined condition being met.

6. The computer-readable storage medium of claim 5 wherein:

the second identification is received in said command, adjacent to the first identification.

7. The computer-readable storage medium of claim 5 further comprising:

instructions to receive with the command, identification of a value of a data type of the at least one argument;

wherein the function is further identified, for use in said instructions to invoke, based at least on the value of the data type.

8. The computer-readable storage medium of claim 5 wherein:

the first structure further comprises a first name of the first aggregation and a second name of the second aggregation; and

the first structure further comprises at least another data type of value to be returned by the aggregation.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 31, 2007
From: SRINIVASAN, ANAND; JAIN, NAMIT; MISHRA, SHAILENDRA KUMAR
To: ORACLE INTERNATIONAL CORPORATION
Reel/Frame 020305/0182 →
Continuity (1)
Related Publication 20090106198A1 · Apr 23, 2009