IP Library › Granted Patent US 11,809,430
Granted Patent B2
US 11,809,430 · App. 16/804,214 · Granted Nov 7, 2023

Efficient stream processing with data aggregations in a sliding window over out-of-order data streams

Inventor: Felix Klaedtke (Heidelberg, DE)
Assignee: NEC CORPORATION
G06F16/24568G06F7/14G06F16/24556H04L63/1416H04L63/20
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,809,430
App. No.
16/804,214
Granted
Nov 7, 2023
Kind
B2
Abstract

A method for processing an out-of-order data stream includes inserting a new data stream element into a segment list according to a timestamp of the new data stream element. It is identified whether there are missing data stream elements between segments in the segment list. The segments which have no missing data stream elements between them are merged. Values of the data stream elements are aggregated using a sliding window over out-of-order data stream elements in the merged segment.

Claims (33)

1. A method for processing an out-of-order data stream, the method comprising:

inserting a new data stream element into a segment list according to a timestamp of the new data stream element;

identifying whether there are missing data stream elements between segments in the segment list;

merging the segments which have no missing data stream elements between them; and

aggregating values of the data stream elements using a sliding window over out-of-order data stream elements in the merged segment,

wherein each of the segments includes a left-most sliding window and a right-most sliding window, wherein the values of the data stream elements are aggregated by moving the right-most sliding window of a first one of the segments to the right and computing data aggregations in each window until a left bound of the right-most sliding window of the first one of the segments matches with a left bound of the left-most sliding window of a second one of the segments, the second one of the segments spanning a time window that is later than the first one of the segments, and wherein the computed data aggregations for each of the windows are output.

2. The method according to claim 1 , further comprising removing data stream elements between a right bound of the left-most sliding window of the first one of the segments and the left bound of the right-most sliding window of the second one of the segments.

3. The method according to claim 2 , wherein a plurality of pairs of segments are merged in parallel.

4. The method according to claim 1 , wherein the segment list is a skip list which stores partial data aggregations, the segments being ordered ascendingly by timestamps of their stream elements, and wherein the new data stream element is inserted into the skip list as a new singleton segment.

5. The method according to claim 4 , wherein the skip list includes a plurality of buckets into which data stream elements of the data stream are insertable in parallel.

6. The method according to claim 1 , further comprising inserting a gap element for an identified missing data stream element.

7. The method according to claim 6 , wherein the gap element has meta-information which includes a timestamp of a singleton interval and a sequence number of the missing data element having the timestamp together with an end marker.

8. The method according claim 1 , further comprising providing a lexicographical ordering of the data stream elements by annotating each data stream element of the data stream from a plurality of data producers with sequence numbers.

9. The method according to claim 8 , further comprising filtering some of the data stream elements out of the data stream and inserting gap elements annotated with the same sequence numbers as the data stream elements which were filtered out.

10. The method according to claim 1 , further comprising inserting a gap element for an identified missing data stream element, the inserted gap element being annotated with meta-information including a timestamp of a time window of the segments, a data producer and a sequence number.

11. The method according to claim 10 , wherein the data producer is a data producer of a first data stream element in the time window, and wherein the sequence number comprises two parts, a first part having a sequence number of the first data stream element and a second part having a counter value of a number of time windows that start at the timestamp.

12. The method according to claim 1 , wherein a tree is stored for each segment in the segment list, wherein the data stream elements of the segments are aggregated using an associative operator from left to right, and wherein the subtrees of the trees of the segments are reused during the aggregation.

13. A system comprising one or more processors which, alone or in combination, are configured to provide for execution of a method for processing an out-of-order data stream, the method comprising:

inserting a new data stream element into a segment list according to a timestamp of the new data stream element;

identifying whether there are missing data stream elements between segments in the segment list;

merging the segments which have no missing data stream elements between them; and

aggregating values of the data stream elements using a sliding window over out-of-order data stream elements in the merged segment,

wherein each of the segments includes a left-most sliding window and a right-most sliding window, wherein the values of the data stream elements are aggregated by moving the right-most sliding window of a first one of the segments to the right and computing data aggregations in each window until a left bound of the right-most sliding window of the first one of the segments matches with a left bound of the left-most sliding window of a second one of the segments, the second one of the segments spanning a time window that is later than the first one of the segments, and wherein the computed data aggregations for each of the windows are output.

14. A tangible, non-transitory computer-readable medium having instructions thereon which, upon being executed by one or more processors, alone or in combination, provide for execution of a method for processing an out-of-order data stream, the method comprising:

inserting a new data stream element into a segment list according to a timestamp of the new data stream element;

identifying whether there are missing data stream elements between segments in the segment list;

merging the segments which have no missing data stream elements between them; and

aggregating values of the data stream elements using a sliding window over out-of-order data stream elements in the merged segment,

wherein each of the segments includes a left-most sliding window and a right-most sliding window, wherein the values of the data stream elements are aggregated by moving the right-most sliding window of a first one of the segments to the right and computing data aggregations in each window until a left bound of the right-most sliding window of the first one of the segments matches with a left bound of the left-most sliding window of a second one of the segments, the second one of the segments spanning a time window that is later than the first one of the segments, and wherein the computed data aggregations for each of the windows are output.

15. The system according to claim 13 , wherein the method further comprises removing data stream elements between a right bound of the left-most sliding window of the first one of the segments and the left bound of the right-most sliding window of the second one of the segments.

16. The system according to claim 15 , wherein a plurality of pairs of segments are merged in parallel.

17. The tangible, non-transitory computer-readable medium according to claim 14 , wherein the method further comprises removing data stream elements between a right bound of the left-most sliding window of the first one of the segments and the left bound of the right-most sliding window of the second one of the segments.

18. The tangible, non-transitory computer-readable medium according to claim 17 , wherein a plurality of pairs of segments are merged in parallel.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 2, 2023
From: NEC LABORATORIES EUROPE GMBH
To: NEC CORPORATION
Reel/Frame 065088/0771 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 2, 2020
From: KLAEDTKE, FELIX
To: NEC LABORATORIES EUROPE GMBH
Reel/Frame 051974/0299 →
Continuity (2)
Provisional Application 62924709 · Oct 23, 2019
Related Publication 20210124746A1 · Apr 29, 2021