US2014040903A1PendingUtilityA1
Queue and operator instance threads to losslessly process online input streams events
Est. expiryJul 31, 2032(~6 yrs left)· nominal 20-yr term from priority
G06F 9/5027G06F 2209/5018
42
PatentIndex Score
0
Cited by
0
References
0
Claims
Abstract
A queue enqueues an online input stream of events arriving at the queue in real-time. An operator instance has one or more threads to losslessly dequeue and process the events from the queue, and to output processing results of the events in a common output stream. The one or more threads are dynamically instantiated and destantiated to maintain an optimal number of the one or more threads while ensuring that none of the events of the online input stream are dropped.
Claims
exact text as granted — not AI-modifiedWe claim:
1 . An apparatus comprising:
a processor; a computer-readable data storage medium; a queue implemented by the processor at the computer-readable data storage medium to enqueue an online input stream of events arriving at the queue in real-time; an operator instance implemented by the processor and having one or more threads to losslessly dequeue and process the events from the queue, and to output processing results of the events in a common output stream; and a control mechanism implemented by the processor to dynamically instantiate and destantiate the one or more threads to maintain an optimal number of the one or more threads while ensuring that none of the events of the online input stream are dropped.
2 . The apparatus of claim 1 , wherein the queue has a static size regardless of a variable arrival rate of the events at the queue.
3 . The apparatus of claim 1 , wherein the events arrive at the queue at a variable arrival rate, such that the control mechanism is to decrease a number of the one or more threads as the variable arrival rate of the events at the queue decreases and is to increase the number of the one or more threads as the variable arrival rate of the events at the queue increases.
4 . The apparatus of claim 1 , wherein the one or more threads dequeue and process the events asynchronously in relation to arrival of the events at the queue.
5 . The apparatus of claim 1 , wherein each thread upon becoming available dequeues a next event from the queue, processes the next event, and outputs a processing result of the next event in the common output stream in coordination with other threads of the one or more threads, such that the processing results within the common output stream are ordered in accordance with an order in which the events are enqueued within the queue,
and wherein the common output stream is common to each thread.
6 . The apparatus of claim 1 , wherein the control mechanism is to monitor a fullness of the queue, is to increase the number of the one or more threads as the fullness of the queue increases, and is to decrease the number of the one or more threads as the fullness of the queue decreases.
7 . The apparatus of claim 1 , wherein the operator instance is a single and only operator instance to dequeue and process the events of the online input stream, and the common output stream is a single and only output stream in which the processing results of the events of the online input stream are output.
8 . A method comprising:
adding to a queue an online input stream of events arriving at the queue in real-time; by each thread of one or more threads of an operator instance,
removing a next event from the queue;
processing the next event removed from the queue; and
outputting a processing result of the next event within an output stream common to the one or more threads,
wherein the events are losslessly added to and removed from the queue.
9 . The method of claim 8 , wherein removing the next event from the queue, processing the next event removed from the queue, and outputting the processing result of the next event within the output stream are performed asynchronously in relation to arrival of the events at the queue.
10 . The method of claim 8 , further comprising:
dynamically instantiating and destantiating the one or more threads to maintain an optimal number of the one or more threads while ensuring that none of the events of the online input stream are dropped, wherein the events arrive at the queue at a variable arrival rate, such that a number of the one or more threads is decreased as the variable arrival rate decreases and is increased as the variable arrival rate increases.
11 . The method of claim 10 , further comprising:
monitoring fullness of the queue, such that the number of the one or more threads is increased as the fullness increases and is decreased as the fullness decreases.
12 . The method of claim 8 , wherein the operator instance is a single and only operator instance to dequeue and process the events of the online input stream, and the common output stream is a single and only output stream in which the processing results of the events of the online input stream are output.
13 . A non-transitory computer-readable data storage medium storing a computer program executable by a processor to perform a method comprising:
adding and removing threads of an operator instance that losslessly dequeue and process events of an online input stream that arrive at and are enqueued within a queue in real-time, to maintain an optimal number of the threads while ensuring that none of the events of the online input stream are dropped, wherein the threads output processing results in a common output stream.
14 . The non-transitory computer-readable data storage medium of claim 13 , wherein the events arrive at the queue at a variable arrival rate,
and wherein adding and removing the threads of the operator instance comprises decreasing a number of the threads as the variable arrival rate decreases and increasing the number of the threads as the variable arrival rate increases.
15 . The non-transitory computer-readable data storage medium of claim 13 , wherein the method further comprises:
monitoring a fullness of the queue, and wherein adding and removing the threads of the operator instance comprises increasing a number of the threads as the fullness increases and decreasing the number of the threads as the fullness decreases.Join the waitlist — get patent alerts
Track US2014040903A1 — get alerts on status changes and closely related new filings.
We store only your email — no account needed. See our privacy policy.