IP Library Granted Patent US 8,954,972
Granted Patent B2
US 8,954,972 · App. 13/553,651 · Granted Feb 10, 2015

Systems and methods for event stream processing

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 8,954,972
App. No.
13/553,651
Granted
Feb 10, 2015
Kind
B2
Abstract

Disclosed are systems and methods for processing events in an event stream using a map-update application. The events may be embodied as a key-attribute pair. An event is processed by one or more instances implementing either a map or an update function. A map function receives an input event from the event stream and publishes one or more events to the event stream. An update function receives an event and updates a corresponding slate and publishes zero or more events. Systems and methods are also disclosed herein for implementing a map-update application in a multithreaded architecture and for handling overloading of a particular thread or node. Systems and methods for providing access to slates updated according to update operations are also disclosed.

Claims (120)

1. A method for performing event processing comprising:

evaluating an event stream;

for each event in the event stream:

mapping the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process;

mapping the event to a target instance, the target instance including one of a map function and an update function;

selecting a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

performing the one of a map function and an update function of the target instance on the event in the selected thread;

wherein selecting the selected thread from one of the primary thread and the secondary thread further comprises:

evaluating if one of the primary thread and the secondary thread is a hotspot thread for the target instance; and

selecting the hotspot thread as the selected thread if one of the primary thread and secondary thread is the hotspot thread for the target instance;

wherein selecting the selected tread from one primary thread and the secondary thread further comprises, if neither of the primary thread and secondary thread is a hotspot thread for the target instance:

evaluating loading of the primary thread and secondary thread; and

if loading of the primary thread is a threshold amount above loading of the secondary thread, selecting the secondary thread as the selected thread, otherwise selecting the primary thread as the selected thread.

2. The method of claim 1 wherein evaluating loading of the primary thread and secondary thread comprises evaluating event queues each associated with one of the primary thread and secondary thread.

3. The method of claim 1 , wherein whichever of the primary thread and secondary thread that is currently processing a prior event mapped to the target instance is the hotspot thread for the target instance.

4. The method of claim 1 , wherein the multithreaded process is a virtual machine.

5. The method of claim 1 , wherein the selected node includes a system with a plurality of processors and wherein the multithreaded process includes a number of threads corresponding to a number of the plurality of processes.

6. A method for performing event processing comprising:

evaluating and event stream;

for each event in the event stream:

mapping the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process:

mapping the event to a target instance, the target instance including one of a map function and an update function;

selecting a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

performing the one map function and an update function of the target instance on the event in the selected thread;

wherein the event is a first event; and

wherein each of the primary thread and secondary thread has a queue associated therewith, the method further comprising:

adding the first event to the queue associated with the selected thread;

retrieving the first event from the queue associated with the selected thread; and

if an opposite thread of the primary and secondary threads that is not the selected thread is processing a second event mapped to the target instance, transferring the first event to the queue associated with the opposite thread.

7. The method of claim 6 , wherein the target instance includes an update function and has a slate associated therewith; and

wherein processing the second event comprises:

performing the update function on the second event and updating the slate according to the update function.

8. The method of claim 7 , wherein performing the update function on the second event and updating the slate according to the update function further comprises locking the slate.

9. A method for performing event processing comprising:

evaluating an event stream;

for each event in the stream:

mapping the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithread process;

mapping the event to a target instance, the target instance including one of a map function and an update function;

selecting a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

performing the one of a map function and an update function of the target instance on the event in the selected thread:

wherein the target instance is a first target instance; and

wherein each of the primary thread and secondary thread has a queue associated therewith, the method further comprising:

adding the event to the queue associated with the selected thread;

retrieving the event from the queue associated with the selected thread; and

if the selected thread is a hotspot for a second target instance, mapping the event to a different thread from the selected thread.

10. A method for performing event processing comprising:

evaluating an event stream;

for each event in the event stream:

mapping the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process:

mapping the event to a target instance, the target instance including one of a map function and an update function:

selecting a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

performing the one of a map function and a update function of the target instance on the event in the selected thread:

wherein the target instance is a first target instance and the event is a first event;

wherein each of the primary thread and secondary thread has a queue associated therewith; and

wherein the method further comprises:

adding the event to the queue associated with the selected thread;

if the selected thread is a hot spot for a second target instance:

retrieving and processing N other events for the second target instance, regardless of queue order, from the queue associated with the selected thread;

following processing of the N other events, retrieving the first event from the queue associated with the selected thread; and

mapping the first event to a different thread from the selected thread.

11. A system for performing event processing comprising one or more processors and one or more memory devices operably coupled to the one or more processors, the one or more memory devices storing executable and operational code effective to cause the one or more processors to:

evaluate an event stream;

for each event in the event stream:

map the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process;

map the event to a target instance, the target instance including one of a map function and an update function;

select a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

perform the one of a map function and an update function of the target instance on the event in the selected thread;

wherein the executable and operational code are further effective to cause the one or more processors to select the selected thread from one of the primary thread and the secondary thread by:

evaluating if one of the primary thread and the secondary thread is a hotspot thread for the target instance; and

selecting the hotspot thread as the selected thread if one of the primary thread and secondary thread is the hotspot thread for the target instance

wherein the executable and operational data are further effective to cause the one or more processors to select the selected thread from one of the primary thread and the secondary thread further by, if neither of primary thread and secondary thread is a hotspot thread for the target instance:

evaluating loading of the primary thread and secondary thread; and

if loading of the primary thread is a threshold amount above loading of the secondary thread, selecting the secondary thread as the selected thread, otherwise selecting the primary thread as the selected thread.

12. The system of claim 11 , wherein evaluating loading of the primary thread and secondary thread comprises evaluating event queues each associated with one of the primary thread and secondary thread.

13. The method of claim 11 , wherein whichever of the primary thread and secondary thread that is currently processing a prior event mapped to the target instance is the hotspot thread for the target instance.

14. The system of claim 11 , wherein the multithreaded process is a virtual machine.

15. The system of claim 11 , wherein the selected node includes a system with a plurality of processors and wherein the multithreaded process includes a number of threads corresponding to a number of the plurality of processes.

16. A system for performing event processing comprising one or more processors and one or more memory devices operably coupled to the one or more processors, the one or more memory devices storing executable and operational code effective to cause the one or more processors to:

evaluate an event stream;

for each event in the event stream:

map the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process;

map the event to a target instance, the target instance including one of a map function and an update function:

select a selected thread from one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned to the target instance; and

perform the one map function and an update function of the target instance on the event in the selected thread;

wherein the event is a first event; and

wherein each of the primary thread and secondary thread has a queue associated therewith, the executable and operational code being further effective to cause the one or more processors to:

add the first event to the queue associated with the selected thread;

retrieve the first event from the queue associated with the selected thread; and

if an opposite thread of the primary and secondary threads that is not the selected thread is processing a second event mapped to the target instance, transfer the first event to the queue associated with the opposite thread.

17. The system of claim 16 , wherein the target instance includes an update function and has a slate associated therewith; and

wherein processing the second event comprises:

performing the update function on the second event and updating the slate according to the update function.

18. The system of claim 17 , wherein performing the update function on the second event and updating the slate according to the update function further comprises locking the slate.

19. A system for performing event processing comprising one or more processors and one or more memory devices operably coupled to the one or more processors, the one or more memory devices storing executable and operational code effective to cause the one or more processors to:

evaluate an event stream;

for each event in the event stream:

map the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process;

map the event to a target instance, the target instance including one of a map function and an update function;

select a selected thread form one of a primary thread and a secondary thread of the multithreaded process, the primary thread and secondary thread being assigned t the target instance; and

perform the one of a map function and an update function of the target instance on the event in the selected thread;

wherein the target instance is a first target instance; and

wherein each of the primary thread and secondary thread has a queue associated therewith; and

wherein the executable and operational code are further effective to cause the one or more processors to:

add the event to the queue associated with the selected thread;

retrieve the event from the queue associated with the selected thread; and

if the selected thread is a hotspot for a second target instance, map the event to a different thread from the selected thread.

20. A system for performing event processing comprising one or more processors and one or more memory devices operably couples to the one or more processors, the one or more memory devices storing executable and operational code effective to cause the one or more processors to:

evaluate an event stream;

for each event in the event stream:

Map the event to at least one selected node of a plurality of worker nodes, each of the nodes executing a multithreaded process;

map the event to a target instance, the target instance including one of a map function and an update function;

select a selected thread from one of a primary thread and a secondary thread of the multithread process, the primary thread and secondary thread being assigned to the target instance; and

perform the one of a map function and an update function of the target assigned on the event in the selected thread; wherein the target instance is a first target instance and the event is a first event;

wherein each of the primary thread and secondary thread has a queue associated therewith; and

wherein the executable and operational code are further effective to cause the one or more processors to:

add the event to the queue associated with the selected thread;

if the selected thread is a hot spot for a second target instance:

retrieve and processing N other events for the second target instance, regardless of queue order, from the queue associated with the selected thread;

following processing of the N other events, retrieve the first event from the queue associated with the selected thread; and

map the first event to a different thread from the selected thread.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 2, 2018
From: WAL-MART STORES, INC.
To: WALMART APOLLO, LLC
Reel/Frame 045817/0115 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 3, 2012
From: LAM, WANG CHEE; LIU, LU; SIRIPURAPU, TARAKA SUBRAHMANYA PRASAD; RAJARAMAN, ANAND; VACHERI, ZOHEB; DOAN, ANHAI
To: WAL-MART STORES, INC.
Reel/Frame 029072/0541 →