IP Library › Granted Patent US 11,645,114
Granted Patent B2
US 11,645,114 · App. 17/575,477 · Granted May 9, 2023

Distributed streaming system supporting real-time sliding windows

Inventors: João Miguel Forte Oliveirinha (Loures, PT); Ana Sofia Leal Gomes (Lisbon, PT); Pedro Cardoso Lessa e Silva (Oporto, PT); Pedro Gustavo Santos Rodrigues Bizarro (Lisbon, PT)
G06F9/4887G06F9/505G06F9/5038G06F9/542G06F11/3072G06F16/211G06F2201/835G06F2201/86
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,645,114
App. No.
17/575,477
Granted
May 9, 2023
Kind
B2
Abstract

In various embodiments, a process for providing a distributed streaming system supporting real-time sliding windows includes receiving a stream of events at a plurality of distributed nodes and routing the events into topic groupings. The process includes using one or more events in at least one of the topic groupings to determine one or more metrics of events with at least one window and an event reservoir including by: tracking, in a volatile memory of the event reservoir, beginning and ending events within the at least one window; and tracking, in a persistent storage of the event reservoir, events associated with tasks assigned to a respective node. The process includes updating the one or more metrics based on one or more previous values of the one or more metrics as a new event is added or an existing event is expired from the at least one window.

Claims (54)

1. A method, comprising:

receiving a stream of events at a plurality of distributed nodes;

routing the events into topic groupings;

using one or more events in at least one of the topic groupings to determine one or more metrics of events with at least one window and an event reservoir including by:

tracking, in a volatile memory of the event reservoir, beginning and ending events within the at least one window; and

tracking, in a persistent storage of the event reservoir, events associated with tasks assigned to a respective node of the plurality of distributed nodes; and

updating the one or more metrics based on one or more previous values of the one or more metrics as a new event is added or an existing event is expired from the at least one window.

2. The method of claim 1 , further comprising:

sliding the at least one window as the new event is added or the existing event is expired, wherein the at least one window includes a real-time sliding window.

3. The method of claim 1 , wherein using the one or more events in at least one of the topic groupings to determine the one or more metrics of events with the at least one window includes updating an execution task plan.

4. The method of claim 1 , wherein using the one or more events in at least one of the topic groupings to determine the one or more metrics of events with the at least one window includes distributing tasks to the plurality of distributed nodes such that computation of tasks is scalable by adding nodes.

5. The method of claim 1 , further comprising:

assigning a task to a node of the plurality of distributed nodes,

meeting a load budget of each node of the plurality of distributed nodes, or

replicating at least a portion of the tasks to nodes in the plurality of distributed nodes.

6. The method of claim 1 , further comprising, using an execution task plan to, for at least one active task in a set of active tasks:

attempt to assign an active task to an active processor; and

in response to a determination that the attempt to assign the active task to the active processor has failed, attempt to assign the active task to a replica processor.

7. The method of claim 6 , further comprising, after assigning the active task, using an execution task plan to, for at least one replica task in a set of replica tasks:

attempt to assign a replica task to a replica processor;

in response to a determination that the attempt to assign the replica task to the replica processor has failed, attempt to assign the active task to a stale processor; and

in response to a determination that the attempt to assign the active task to the stale processor has failed, assign the active task a processor meeting at least one criterion.

8. The method of claim 7 , wherein the assignment of the active task and the assignment of the replica task are performed during a recovery operation.

9. The method of claim 1 , wherein the at least one window includes at least one of: a real-time sliding window, a delayed window, an infinite window, or a tumbling window.

10. The method of claim 1 , further comprising determining a task encapsulating calculation of all metrics associated with a specified partition.

11. The method of claim 10 , wherein the specified partition is a unit of computation such that events are processed according to their respective partitions.

12. The method of claim 10 , wherein the task is at least one of: an active task assigned to a node of the plurality of nodes or a replica task for which the node is a backup processor.

13. The method of claim 1 , wherein a replication factor associated with how many replica tasks to create is based on a threshold of tolerable failures.

14. The method of claim 1 , wherein:

events of the stream of events are organized into chunks including member events ordered by timestamp; and

at least one adjacent chunk is eagerly loaded in response to the loading of one of the chunks into a volatile memory.

15. The method of claim 1 , wherein the event reservoir includes a schema registry describing fields associated with events and the schema registry is updated in response to schema changes within the stream of events.

16. The method of claim 1 , wherein:

a metric state store is organized as a key-value store;

a key in the key-value store represents a metric entity in a task plan; and

a number of keys accessed per event matches a number of roots in the task plan.

17. The method of claim 16 , further comprising providing checkpointing, wherein checkpoint triggers are synchronized between the event reservoir and the metric state store.

18. The method of claim 1 , wherein routing the events into the topic groupings includes replicating a given event as many times as a number of partitioners for the stream.

19. A system, comprising:

an event reservoir including a volatile memory and a persistent storage;

one or more processors included in a plurality of distributed nodes, at least a portion of is the one or more processors configured to:

receive a stream of events;

route the events into topic groupings;

use one or more events in at least one of the topic groupings to determine one or more metrics of events with at least one window and an event reservoir including by:

tracking, in the volatile memory of the event reservoir, beginning and ending events within the at least one window; and

tracking, in the persistent storage of the event reservoir, events associated with assigned tasks; and

update the one or more metrics based on one or more previous values of the one or more metrics as a new event is added or an existing event is expired from the at least one window.

20. A computer program product embodied in a non-transitory computer readable medium and comprising computer instructions for:

receiving a stream of events at a plurality of distributed nodes;

routing the events into topic groupings;

using one or more events in at least one of the topic groupings to determine one or more metrics of events with at least one window and an event reservoir including by:

tracking in a volatile memory of the event reservoir, beginning and ending events within the at least one window; and

tracking, in a persistent storage of the event reservoir, events associated with tasks of a respective node of the plurality of distributed nodes; and

updating the one or more metrics based on one or more previous values of the one or to more metrics as a new event is added or an existing event is expired from the at least one window.

Priority Claims (2)
WO 21174843 · May 19, 2021 · international
PT 117243 · May 19, 2021 · national
Continuity (3)
Continuation 17356310 · Jun 23, 2021
Provisional Application 63066035 · Aug 14, 2020
Related Publication 20220138006A1 · May 5, 2022
Cited By (1)
US 12,748,674