IP Library › Granted Patent US 11,487,764
Granted Patent B2
US 11,487,764 · App. 16/827,122 · Granted Nov 1, 2022

System and method for stream processing

Inventors: Radu Tudoran (Munich, DE); Stefano Bortoli (Munich, DE); Xing Zhu (Shanghai, CN); Goetz Brasche (Munich, DE); Cristian Axenie (Munich, DE)
Assignee: Huawei Cloud Computing Technologies Co., Ltd.
G06F16/24568G06F16/24552
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 11,487,764
App. No.
16/827,122
Granted
Nov 1, 2022
Kind
B2
Abstract

An input stream of events is processed to obtain an output stream of events. Consecutive events are selected from the input stream using a sliding window to obtain sliding window events, then a function is applied thereto to obtain an output result value. Operations of: outputting the output result value in the output stream; splitting the sliding window events into filter-complying events and pending events; applying the function on the pending events to obtain preliminary value(s); selecting, from the input stream, a second plurality of events; adding the second plurality of events to the sliding window events; removing, from the sliding window events, the filter-complying events to obtain a new set of sliding window events; and applying the function to the second plurality of events and the preliminary value(s) to obtain a new output result value, are then iteratively performed.

Claims (54)

1. A method for processing an input stream of events to obtain an output stream of events, wherein each of the events of the input stream has an event value, the method comprising:

selecting, from the input stream of events, a plurality of consecutive events using a sliding window to obtain sliding window events;

applying a function to the event values of the sliding window events to obtain an output result value; and

in each of a plurality of iterations:

outputting the output result value in the output stream of events;

splitting the sliding window events into a set of complying events, satisfying at least one filter test, and a set of remaining pending events;

storing the set of complying events in a cache memory, the cache memory being a random access memory, and storing the set of remaining pending events in a non-volatile memory;

applying the function to the event values of the set of remaining pending events to obtain at least one preliminary value and storing the at least one preliminary value in the cache memory;

selecting, from the input stream of events, a second plurality of events, wherein the second plurality of events comprises new events that are newly received at the input stream of events;

adding the second plurality of events to the sliding window events;

removing, from the sliding window events, the set of complying events to obtain a new set of sliding window events;

retrieving the at least one preliminary value from the cache memory;

applying the function to the event values of the second plurality of events to obtain a head result; and

applying the function to the head result and the at least one preliminary value to obtain a new output result value.

2. The method according to claim 1 , wherein the at least one preliminary value, after the applying of the function on the set of remaining pending events, comprises at least one of:

an average value of a plurality of event values of a plurality of events of the input stream of events,

a minimum value of the plurality of event values,

a maximum value of the plurality of event values,

an amount of values in the plurality of event values,

an amount of distinct values in the plurality of event values,

a sum of the plurality of event values,

a median value of the plurality of event values,

a quartile value of the plurality of event values,

a standard deviation value of the plurality of event values, or

a variance value of the plurality of event values.

3. The method according to claim 1 , wherein the step of applying the function on the set of remaining pending events comprises:

splitting the plurality of remaining pending events into a plurality of buckets according to a second filter test;

applying the function on each bucket of the plurality of buckets to obtain a corresponding plurality of preliminary values; and

storing the plurality of preliminary values in the cache memory.

4. The method according to claim 3 , wherein the plurality of preliminary values comprise a plurality of minimum or maximum values of a plurality of event values of the plurality of events in one of the plurality of buckets of the pending events.

5. The method according to claim 3 , wherein the plurality of preliminary values comprise a plurality of bucket count values, each bucket count value, of the bucket count values, counting an amount of events in one of the plurality of buckets of the pending events.

6. The method according to claim 3 , wherein all event values in a first one of the plurality of buckets of the pending events succeed all event values in a second one of the plurality of buckets of the pending events according to an identified ordering function.

7. The method according to claim 1 ,

wherein each event of the input stream of events has a time value selected from a group consisting of a time of arrival, a time of creation, and a time of occurrence of the event, and

wherein at least one event of the plurality of complying events has a time value that is earlier than a time value of any event of the set of remaining pending events.

8. The method according to claim 1 , wherein the at least one filter test comprises:

comparing a time of an event to a certain time range relative to a present time; or

comparing a value of an event to one or more threshold values.

9. A system for processing an input stream of events to obtain an output stream of events, wherein each of the events of the input stream has an event value the system comprising a processor that is configured to:

select, from the input stream of events, a plurality of consecutive events using a sliding window to obtain sliding window events;

apply a function to the event values of the sliding window events to obtain an output result value; and

in each of a plurality of iterations:

output the output result value in the output stream of events;

split the sliding window events into a set of complying events, satisfying at least one filter test, and a set of remaining pending events;

store the set of complying events in a cache memory, the cache memory being a random access memory, and store the set of remaining pending events in a non-volatile memory;

apply the function to the event values of the set of remaining pending events to obtain at least one preliminary value and store the at least one preliminary value in the cache memory;

select, from the input stream of events, a second plurality of events, wherein the second plurality of events comprises new events that are newly received at the input stream of events;

add the second plurality of events to the sliding window events;

remove, from the sliding window events, the set of complying events to obtain a new set of sliding window events;

retrieve the at least one preliminary value from the cache memory;

apply the function to the event values of the second plurality of events to obtain a head result, and

apply the function to the head result and the at least one preliminary value, to obtain a new output result value.

10. The system according to claim 9 , wherein the non-volatile memory comprises one of a hard disk electrically connected to the processor, a network memory connected to the processor via a network interface, a database, a local file system, a distributed file system, or a cloud storage.

11. A non-transitory computer readable medium comprising program code configured to perform the method according to claim 1 upon the computer program being executed on a computer.

Assignments (3)
CORRECTIVE ASSIGNMENT TO CORRECT THE THE EXECUTION DATE OF THE INVENTOR 3RD INVENTOR. PREVIOUSLY RECORDED AT REEL: 060942 FRAME: 0861. ASSIGNOR(S) HEREBY CONFIRMS THE ASSIGNMENT. Recorded Sep 21, 2022
From: TUDORAN, RADU; BORTOLI, STEFANO; ZHU, XING; BRASCHE, GOETZ; AXENIE, CRISTIAN
To: HUAWEI TECHNOLOGIES CO., LTD.
Reel/Frame 061498/0635 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 30, 2022
From: TUDORAN, RADU; BORTOLI, STEFANO; ZHU, XING; BRASCHE, GOETZ; AXENIE, CRISTIAN
To: HUAWEI TECHNOLOGIES CO., LTD.
Reel/Frame 060942/0861 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 1, 2022
From: HUAWEI TECHNOLOGIES CO., LTD.
To: HUAWEI CLOUD COMPUTING TECHNOLOGIES CO., LTD.
Reel/Frame 059267/0088 →
Continuity (2)
Continuation PCTEP2017073956 · Sep 21, 2017
Related Publication 20200285646A1 · Sep 10, 2020