【问题标题】:Message queue with time to live有时间的消息队列
【发布时间】:2016-10-15 00:39:30
【问题描述】:

我寻找一个队列,该队列最多可存储 N 元素一段时间(即 10 秒),或者如果已满则应处置最旧的值。

我在 Apache Collections (CircularFifoQueueJavaDoc) 中发现了一个类似的队列,它忽略了生存时间。对于我的小项目来说,一个完整的消息代理似乎是一种开销。

你介意给我一个提示我应该如何实现吗?轮询元素时是否应过滤掉旧值?

【问题讨论】:

    标签: java collections message-queue java.util.concurrent


    【解决方案1】:

    我习惯于跟随队列实现。该代码主要基于 Apaches CircularFifoQueue,并且仅经过弱测试。此外,该实现不是线程安全的并且不可序列化

    如果您发现错误,请发表评论。

    import java.util.AbstractCollection;
    import java.util.Arrays;
    import java.util.Collection;
    import java.util.Iterator;
    import java.util.NoSuchElementException;
    import java.util.Queue;
    import java.util.concurrent.TimeUnit;
    
    /**
     * TimedQueue is a first-in first-out queue with a fixed size that
     * replaces its oldest element if full.
     * <p>
     * The removal order of a {@link TimedQueue} is based on the
     * insertion order; elements are removed in the same order in which they
     * were added.  The iteration order is the same as the removal order.
     * <p>
     * The {@link #add(Object)}, {@link #remove()}, {@link #peek()}, {@link #poll},
     * {@link #offer(Object)} operations all perform in constant time.
     * All other operations perform in linear time or worse.
     * <p>
     * This queue prevents null objects from being added and it is not thread-safe and not serializable.
     * 
     * The majority of this source code was copied from Apaches {@link CircularFifoQueue}: http://commons.apache.org/proper/commons-collections/apidocs/org/apache/commons/collections4/queue/CircularFifoQueue.html
     *
     * @version 1.0
     */
    public class TimedQueue<E> extends AbstractCollection<E>
    implements Queue<E> {
    
    /** Underlying storage array. */
    private Item<E>[] elements;
    
    /** Array index of first (oldest) queue element. */
    private int start = 0;
    
    /**
     * Index mod maxElements of the array position following the last queue
     * element.  Queue elements start at elements[start] and "wrap around"
     * elements[maxElements-1], ending at elements[decrement(end)].
     * For example, elements = {c,a,b}, start=1, end=1 corresponds to
     * the queue [a,b,c].
     */
    private transient int end = 0;
    
    /** Flag to indicate if the queue is currently full. */
    private transient boolean full = false;
    
    /** Capacity of the queue. */
    private final int maxElements;
    
    private TimeUnit unit;
    private int delay;
    
    /**
     * Constructor that creates a queue with the default size of 32.
     */
    public TimedQueue() {
        this(32);
    }
    
    /**
     * Constructor that creates a queue with the specified size.
     *
     * @param size  the size of the queue (cannot be changed)
     * @throws IllegalArgumentException  if the size is &lt; 1
     */
    public TimedQueue(final int size) {
        this(size, 3, TimeUnit.SECONDS);
    }
    
    @SuppressWarnings("unchecked")
    public TimedQueue(final int size, int delay, TimeUnit unit) {
        if (size <= 0) {
            throw new IllegalArgumentException("The size must be greater than 0");
        }
        elements = new Item[size];
        maxElements = elements.length;
        this.unit = unit;
        this.delay = delay;
    }
    
    /**
     * Constructor that creates a queue from the specified collection.
     * The collection size also sets the queue size.
     *
     * @param coll  the collection to copy into the queue, may not be null
     * @throws NullPointerException if the collection is null
     */
    public TimedQueue(final Collection<? extends E> coll) {
        this(coll.size());
        addAll(coll);
    }
    
    /**
     * Returns the number of elements stored in the queue.
     *
     * @return this queue's size
     */
    @Override
    public int size() {
        int size = 0;
    
        for(int i = 0; i < elements.length; i++) {
            if(validElement(i) != null) {
                size++;
            }
        }
    
        return size;
    }
    
    /**
     * Returns true if this queue is empty; false otherwise.
     *
     * @return true if this queue is empty
     */
    @Override
    public boolean isEmpty() {
        return size() == 0;
    }
    
    private boolean isAtFullCapacity() {
        return size() == maxElements;
    }
    
    /**
     * Clears this queue.
     */
    @Override
    public void clear() {
        full = false;
        start = 0;
        end = 0;
        Arrays.fill(elements, null);
    }
    
    /**
     * Adds the given element to this queue. If the queue is full, the least recently added
     * element is discarded so that a new element can be inserted.
     *
     * @param element  the element to add
     * @return true, always
     * @throws NullPointerException  if the given element is null
     */
    @Override
    public boolean add(final E element) {
        if (null == element) {
            throw new NullPointerException("Attempted to add null object to queue");
        }
    
        if (isAtFullCapacity()) {
            remove();
        }
    
        elements[end++] = new Item<E>(element);
    
        if (end >= maxElements) {
            end = 0;
        }
    
        if (end == start) {
            full = true;
        }
    
        return true;
    }
    
    /**
     * Returns the element at the specified position in this queue.
     *
     * @param index the position of the element in the queue
     * @return the element at position {@code index}
     * @throws NoSuchElementException if the requested position is outside the range [0, size)
     */
    public E get(final int index) {
        final int sz = size();
        if (sz == 0) {
            throw new NoSuchElementException(
                    String.format("The specified index (%1$d) is outside the available range because the queue is empty.", Integer.valueOf(index)));
        }
        if (index < 0 || index >= sz) {
            throw new NoSuchElementException(
                    String.format("The specified index (%1$d) is outside the available range [0, %2$d]",
                                  Integer.valueOf(index), Integer.valueOf(sz-1)));
        }
    
        final int idx = (start + index) % maxElements;
        return validElement(idx);
    }
    
    private E validElement(int idx) {
        if(elements[idx] == null){
            return null;
        }
        long diff = System.currentTimeMillis() - elements[idx].getCreationTime();
    
        if(diff < unit.toMillis(delay)) {
            return (E) elements[idx].getValue();
        } else {
            elements[idx] = null;
            return null;
        }
    }
    
    //-----------------------------------------------------------------------
    
    /**
     * Adds the given element to this queue. If the queue is full, the least recently added
     * element is discarded so that a new element can be inserted.
     *
     * @param element  the element to add
     * @return true, always
     * @throws NullPointerException  if the given element is null
     */
    public boolean offer(E element) {
        return add(element);
    }
    
    public E poll() {
        if (isEmpty()) {
            return null;
        }
        return remove();
    }
    
    public E element() {
        if (isEmpty()) {
            throw new NoSuchElementException("queue is empty");
        }
        return peek();
    }
    
    public E peek() {
        if (isEmpty()) {
            return null;
        }
        return (E) elements[start].getValue();
    }
    
    public E remove() {
        if (isEmpty()) {
            throw new NoSuchElementException("queue is empty");
        }
    
        final E element = validElement(start);
        if (null != element) {
            elements[start++] = null;
    
            if (start >= maxElements) {
                start = 0;
            }
            full = false;
        }
        return element;
    }
    
    /**
     * Increments the internal index.
     *
     * @param index  the index to increment
     * @return the updated index
     */
    private int increment(int index) {
        index++;
        if (index >= maxElements) {
            index = 0;
        }
        return index;
    }
    
    /**
     * Decrements the internal index.
     *
     * @param index  the index to decrement
     * @return the updated index
     */
    private int decrement(int index) {
        index--;
        if (index < 0) {
            index = maxElements - 1;
        }
        return index;
    }
    
    /**
     * Returns an iterator over this queue's elements.
     *
     * @return an iterator over this queue's elements
     */
    @Override
    public Iterator<E> iterator() {
        return new Iterator<E>() {
    
            private int index = start;
            private int lastReturnedIndex = -1;
            private boolean isFirst = full;
    
            public boolean hasNext() {
                return (isFirst || index != end) && size() > 0;
            }
    
            public E next() {
                if (!hasNext()) {
                    throw new NoSuchElementException();
                }
                isFirst = false;
                lastReturnedIndex = index;
                index = increment(index);
                if(validElement(lastReturnedIndex) == null) {
                    return next();
                } else {
                    return validElement(lastReturnedIndex);
                }
            }
    
            public void remove() {
                if (lastReturnedIndex == -1) {
                    throw new IllegalStateException();
                }
    
                // First element can be removed quickly
                if (lastReturnedIndex == start) {
                    TimedQueue.this.remove();
                    lastReturnedIndex = -1;
                    return;
                }
    
                int pos = lastReturnedIndex + 1;
                if (start < lastReturnedIndex && pos < end) {
                    // shift in one part
                    System.arraycopy(elements, pos, elements, lastReturnedIndex, end - pos);
                } else {
                    // Other elements require us to shift the subsequent elements
                    while (pos != end) {
                        if (pos >= maxElements) {
                            elements[pos - 1] = elements[0];
                            pos = 0;
                        } else {
                            elements[decrement(pos)] = elements[pos];
                            pos = increment(pos);
                        }
                    }
                }
    
                lastReturnedIndex = -1;
                end = decrement(end);
                elements[end] = null;
                full = false;
                index = decrement(index);
            }
    
        };
    }
    
    private static final class Item<E> {
        private long creationTime;
        private E in;
    
        public Item(E in) {
            this.in = in;
            creationTime = System.currentTimeMillis();
        }
    
        public E getValue() {
            return in;
        }
    
        public long getCreationTime() {
            return creationTime;
        }
    }
    }
    

    【讨论】:

      【解决方案2】:

      我认为java.util.LinkedHashMap 是适合您的解决方案。它有一个removeEldest() 方法,只要将条目放入地图中就会调用该方法。您可以覆盖它以指示是否应删除最旧的条目。

      JavaDoc 给出了一个很好的例子:

       private static final int MAX_ENTRIES = 100;
      
       protected boolean removeEldestEntry(Map.Entry eldest) {
          return size() > MAX_ENTRIES;
       }
      

      如果地图有超过 100 个条目,这将删除最旧的条目。

      10 秒后主动移除项目需要单独的线程来检查年龄并移除旧项目。从你的描述来看,我猜这不是你想要的。

      【讨论】:

      • 这个答案应该是评论。
      • 这看起来是个好主意。但是,我想删除超过 10 秒的项目。是否已经实现了这种“基于时间的队列”?
      【解决方案3】:

      有一个名为LinkedHashMap 的类,它有一个特殊的方法来删除过时的数据。来自文档:

      protected boolean removeEldestEntry(Map.Entry eldest)
      如果此映射应删除其最旧的条目,则返回 true。

      只要将任何内容添加到列表(队列)中,就会调用方法removeEldestEntry。如果它返回true,则删除最旧的条目以为新条目腾出空间,否则不会删除任何内容。您可以添加自己的实现来检查最旧条目上的时间戳,如果它早于过期阈值(例如 10 秒),则返回 true。所以你的实现可能看起来像这样:

      protected boolean removeEldestEntry(Map.Entry eldest) {
          long currTimeMillis = System.currentTimeMillis();
          long entryTimeMillis = eldest.getValue().getTimeCreation();
      
          return (currTimeMillis - entryTimeMillis) > (1000*10*60); 
      }
      

      【讨论】:

      • 我会检查并及时通知您!
      • 这个答案很棒!
      • @Markus 这对您来说是个问题吗?旧值最终将全部删除。我相信我的回答会很好地满足您的需求。
      • 如果您愿意,只需限制队列的大小,如果有新条目进入,请使用我的解决方案踢掉任何陈旧条目。我相信您可以使当前界面适用于您的用例。
      • @Markus 我支持任何可以让您解决问题的方法。如果我能提供进一步的帮助,请告诉我。
      猜你喜欢
      • 2022-06-14
      • 2016-03-19
      • 1970-01-01
      • 2013-03-09
      • 1970-01-01
      • 2011-05-23
      • 1970-01-01
      • 2013-11-08
      • 2020-11-27
      相关资源
      最近更新 更多