Keep track of top n values in a moving window of size > n?

Viewed 97

Suppose we have a data-stream feeding in once per minute and we want to keep a running track of the top 5 values over the previous 10 minutes. Intuitively there should be some queue solution to this but I'm struggling to find an elegant way due to the fact that an element can be popped for two different reasons (it look place 11 minutes ago or <10 minutes ago but 5 higher values have been seen since)

Any suggestions would be appreciated!

2 Answers

You will not want to implement a Queue because the conveniences of queues are lost as soon as you need to take elements off not just the top (as is the case of smaller-than elements). Here's a custom class implementation that does what you are asking for -- the trick is to check the values of the lists before adding values. This is very simple, but I hope it gives you an idea:

from datetime import datetime, timedelta

class TopLatestN():
    def __init__(self, max_size: int=5, timeframe_m: int=10):
        self.__values = []
        self.times = []  # Matching queue of times items are added
        self.max_size = max_size
        self.timeframe_m = timeframe_m

    def add_value(self, value):
        # Remove values that are too old
        self.__check_times()
        
        if len(self.__values) < self.max_size:
            # Add the element right away if we aren't 'at capacity'
            self.__values.append(value)
            self.times.append(datetime.now())
        else:
            # Add the value only if it's large enough
            new_index = self.__check_values()
            if new_index is not None:
                self.__values[new_index] = value
                self.times[new_index] = datetime.now()

    def get_values(self):
        self.__check_times()
        
        return self.__values

    def __check_times(self):
        # Remove values/times that are too old
        current_time = datetime.now()
        
        # Get matching list of values to keep
        keep = []
        for time in self.times:
            keep.append(time + timedelta(minutes=self.timeframe_m) >= current_time)

        # Replace values for all too-told times
        self.__values = [val for i, val in enumerate(self.__values) if keep[i]]
        self.times = [time for i, time in enumerate(self.times) if keep[i]]

    def __check_values(self, value):
        # Get index to store value at IF it is large enough - else return None
        if any(value > val for val in self.__values):
            return self.__values.index(min(self.__values)) # Replace smallest value

        return None

Don't overlook a simple thing here - you may be surprised at how well it works. You have two orders, so maintain two sequences in sorted order, bytime ordered by time and byvalue ordered by value. Store 2-tuple (value, timestamp) pairs in each. Of course you need to keep them in synch.

Because byvalue is always maintained in order sorted by value, you can, at any time, look at the top n, bottom n, middle n, or any other kind of order statistic you want.

Assuming your timestamps (whatever that may mean to you) only increase over time, "sorting" by time is trivial: use a collections.deque and push new records on one end (say, the right) and discard from the other end. Use a plain list for byvalue. To expire old records, then:

oldest_to_retain = whatever form of timestamp you use
while bytime and bytime[0][1] < oldest_to_retain:
    t = bytime.popleft() # discard expired record
    # and remove it from the other seq too
    i = bisect.bisect_left(byvalue, t)
    assert byvalue[i] == t
    del byvalue[i]

To insert an incoming value,

t = (the_new_value, current_timestamp)
assert not bytime or bytime[-1][1] <= current_timestamp
bytime.append(t)
bisect.insort(byvalue, t)

Now with some experience, people balk at this idea because these statements have O() behavior linear in len(byvalue):

del byvalue[i]
bisect.insort(byvalue, t)

(And the other statements have O(1) or O(log(N)) behavior.)

With more experience, they get over that ;-) Those occur "at C speed", and unless byvalue grows to several hundreds of elements it's generally faster - and more space-efficient - than fancy tree structures, even if they're coded in optimized C.

If byvalue does grow large, it's an easy change to switch byvalue to use a SortedList from the widely used sortedcontainers package. Then no statement is worse than about O(log(N)). Your part of the code remains just as simple, flexible, and easy to reason about.

Related