001/**
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *      http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018
019package org.apache.oozie.util;
020
021import java.util.AbstractQueue;
022import java.util.ArrayList;
023import java.util.Arrays;
024import java.util.Collection;
025import java.util.ConcurrentModificationException;
026import java.util.Iterator;
027import java.util.List;
028import java.util.concurrent.BlockingQueue;
029import java.util.concurrent.DelayQueue;
030import java.util.concurrent.Delayed;
031import java.util.concurrent.FutureTask;
032import java.util.concurrent.TimeUnit;
033import java.util.concurrent.atomic.AtomicInteger;
034import java.util.concurrent.locks.ReentrantLock;
035
036/**
037 * A Queue implementation that support queuing elements into the future and priority queuing.
038 * <p/>
039 * The {@link PriorityDelayQueue} avoids starvation by raising elements priority as they age.
040 * <p/>
041 * To support queuing elements into the future, the JDK <code>DelayQueue</code> is used.
042 * <p/>
043 * To support priority queuing, an array of <code>DelayQueue</code> sub-queues is used. Elements are consumed from the
044 * higher priority sub-queues first. From a sub-queue, elements are available based on their age.
045 * <p/>
046 * To avoid starvation, there is is maximum wait time for an an element in a sub-queue, after the maximum wait time has
047 * elapsed, the element is promoted to the next higher priority sub-queue. Eventually it will reach the maximum priority
048 * sub-queue and it will be consumed when it is the oldest element in the that sub-queue.
049 * <p/>
050 * Every time an element is promoted to a higher priority sub-queue, a new maximum wait time applies.
051 * <p/>
052 * This class does not use a separate thread for anti-starvation check, instead, the check is performed on polling and
053 * seeking operations. This check is performed, the most every 1/2 second.
054 */
055public class PriorityDelayQueue<E> extends AbstractQueue<PriorityDelayQueue.QueueElement<E>>
056        implements BlockingQueue<PriorityDelayQueue.QueueElement<E>> {
057
058    /**
059     * Element wrapper required by the queue.
060     * <p/>
061     * This wrapper keeps track of the priority and the age of a queue element.
062     */
063    public static class QueueElement<E> extends FutureTask<E> implements Delayed {
064        private XCallable<E> element;
065        private int priority;
066        private long baseTime;
067        boolean inQueue;
068
069        /**
070         * Create an Element wrapper.
071         *
072         * @param element element.
073         * @param priority priority of the element.
074         * @param delay delay of the element.
075         * @param unit time unit of the delay.
076         *
077         * @throws IllegalArgumentException if the element is <tt>NULL</tt>, the priority is negative or if the delay is
078         * negative.
079         */
080        public QueueElement(XCallable<E> element, int priority, long delay, TimeUnit unit) {
081            super(element);
082            if (element == null) {
083                throw new IllegalArgumentException("element cannot be null");
084            }
085            if (priority < 0) {
086                throw new IllegalArgumentException("priority cannot be negative, [" + element + "]");
087            }
088            if (delay < 0) {
089                throw new IllegalArgumentException("delay cannot be negative");
090            }
091            this.element = element;
092            this.priority = priority;
093            setDelay(delay, unit);
094        }
095
096        /**
097         * Return the element from the wrapper.
098         *
099         * @return the element.
100         */
101        public XCallable<E> getElement() {
102            return element;
103        }
104
105        /**
106         * Return the priority of the element.
107         *
108         * @return the priority of the element.
109         */
110        public int getPriority() {
111            return priority;
112        }
113
114        /**
115         * Set the delay of the element.
116         *
117         * @param delay delay of the element.
118         * @param unit time unit of the delay.
119         */
120        public void setDelay(long delay, TimeUnit unit) {
121            baseTime = System.currentTimeMillis() + unit.toMillis(delay);
122        }
123
124        /**
125         * Return the delay of the element.
126         *
127         * @param unit time unit of the delay.
128         *
129         * @return the delay in the specified time unit.
130         */
131        public long getDelay(TimeUnit unit) {
132            return unit.convert(baseTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
133        }
134
135        /**
136         * Compare the age of this wrapper element with another. The priority is not used for the comparision.
137         *
138         * @param o the other wrapper element to compare with.
139         *
140         * @return less than zero if this wrapper is older, zero if both wrapper elements have the same age, greater
141         *         than zero if the parameter wrapper element is older.
142         */
143        public int compareTo(Delayed o) {
144            long diff = (getDelay(TimeUnit.MILLISECONDS) - o.getDelay(TimeUnit.MILLISECONDS));
145            if(diff > 0) {
146                return 1;
147            } else if(diff < 0) {
148                return -1;
149            } else {
150                return 0;
151            }
152        }
153
154        /**
155         * Return the string representation of the wrapper element.
156         *
157         * @return the string representation of the wrapper element.
158         */
159        @Override
160        public String toString() {
161            StringBuilder sb = new StringBuilder();
162            sb.append("[").append(element).append("] priority=").append(priority).append(" delay=").
163                    append(getDelay(TimeUnit.MILLISECONDS));
164            return sb.toString();
165        }
166
167    }
168
169    /**
170     * Frequency, in milliseconds, of the anti-starvation check.
171     */
172    public static final long ANTI_STARVATION_INTERVAL = 500;
173
174    protected int priorities;
175    protected DelayQueue<QueueElement<E>>[] queues;
176    protected transient final ReentrantLock lock = new ReentrantLock();
177    private transient long lastAntiStarvationCheck = 0;
178    private long maxWait;
179    private int maxSize;
180    protected AtomicInteger currentSize;
181
182    /**
183     * Create a <code>PriorityDelayQueue</code>.
184     *
185     * @param priorities number of priorities the queue will support.
186     * @param maxWait max wait time for elements before they are promoted to the next higher priority.
187     * @param unit time unit of the max wait time.
188     * @param maxSize maximum size of the queue, -1 means unbounded.
189     */
190    @SuppressWarnings("unchecked")
191    public PriorityDelayQueue(int priorities, long maxWait, TimeUnit unit, int maxSize) {
192        if (priorities < 1) {
193            throw new IllegalArgumentException("priorities must be 1 or more");
194        }
195        if (maxWait < 0) {
196            throw new IllegalArgumentException("maxWait must be greater than 0");
197        }
198        if (maxSize < -1 || maxSize == 0) {
199            throw new IllegalArgumentException("maxSize must be -1 or greater than 0");
200        }
201        this.priorities = priorities;
202        queues = new DelayQueue[priorities];
203        for (int i = 0; i < priorities; i++) {
204            queues[i] = new DelayQueue<QueueElement<E>>();
205        }
206        this.maxWait = unit.toMillis(maxWait);
207        this.maxSize = maxSize;
208        if (maxSize != -1) {
209            currentSize = new AtomicInteger();
210        }
211    }
212
213    /**
214     * Return number of priorities the queue supports.
215     *
216     * @return number of priorities the queue supports.
217     */
218    public int getPriorities() {
219        return priorities;
220    }
221
222    /**
223     * Return the max wait time for elements before they are promoted to the next higher priority.
224     *
225     * @param unit time unit of the max wait time.
226     *
227     * @return the max wait time in the specified time unit.
228     */
229    public long getMaxWait(TimeUnit unit) {
230        return unit.convert(maxWait, TimeUnit.MILLISECONDS);
231    }
232
233    /**
234     * Return the maximum queue size.
235     *
236     * @return the maximum queue size. If <code>-1</code> the queue is unbounded.
237     */
238    public long getMaxSize() {
239        return maxSize;
240    }
241
242    /**
243     * Return an iterator over all the {@link QueueElement} elements (both expired and unexpired) in this queue. The
244     * iterator does not return the elements in any particular order.  The returned <tt>Iterator</tt> is a "weakly
245     * consistent" iterator that will never throw {@link ConcurrentModificationException}, and guarantees to traverse
246     * elements as they existed upon construction of the iterator, and may (but is not guaranteed to) reflect any
247     * modifications subsequent to construction.
248     *
249     * @return an iterator over the {@link QueueElement} elements in this queue.
250     */
251    @Override
252    @SuppressWarnings("unchecked")
253    public Iterator<QueueElement<E>> iterator() {
254        QueueElement[][] queueElements = new QueueElement[queues.length][];
255        lock.lock();
256        try {
257            for (int i = 0; i < queues.length; i++) {
258                queueElements[i] = queues[i].toArray(new QueueElement[0]);
259            }
260        }
261        finally {
262            lock.unlock();
263        }
264        List<QueueElement<E>> list = new ArrayList<QueueElement<E>>();
265        for (QueueElement[] elements : queueElements) {
266            list.addAll(Arrays.asList((QueueElement<E>[]) elements));
267        }
268        return list.iterator();
269    }
270
271    /**
272     * Return the number of elements in the queue.
273     *
274     * @return the number of elements in the queue.
275     */
276    @Override
277    public int size() {
278        int size = 0;
279        for (DelayQueue<QueueElement<E>> queue : queues) {
280            size += queue.size();
281        }
282        return size;
283    }
284
285    /**
286     * Return the number of elements on each priority sub-queue.
287     *
288     * @return the number of elements on each priority sub-queue.
289     */
290    public int[] sizes() {
291        int[] sizes = new int[queues.length];
292        for (int i = 0; i < queues.length; i++) {
293            sizes[i] = queues[i].size();
294        }
295        return sizes;
296    }
297
298    /**
299     * Inserts the specified element into this queue if it is possible to do
300     * so immediately without violating capacity restrictions, returning
301     * <tt>true</tt> upon success and throwing an
302     * <tt>IllegalStateException</tt> if no space is currently available.
303     * When using a capacity-restricted queue, it is generally preferable to
304     * use {@link #offer(Object) offer}.
305     *
306     * @param queueElement the {@link QueueElement} element to add.
307     * @return <tt>true</tt> (as specified by {@link Collection#add})
308     * @throws IllegalStateException if the element cannot be added at this
309     *         time due to capacity restrictions
310     * @throws ClassCastException if the class of the specified element
311     *         prevents it from being added to this queue
312     * @throws NullPointerException if the specified element is null
313     * @throws IllegalArgumentException if some property of the specified
314     *         element prevents it from being added to this queue
315     */
316    @Override
317    public boolean add(QueueElement<E> queueElement) {
318        return offer(queueElement, false);
319    }
320
321    /**
322     * Insert the specified {@link QueueElement} element into the queue.
323     *
324     * @param queueElement the {@link QueueElement} element to add.
325     * @param ignoreSize if the queue is bound to a maximum size and the maximum size is reached, this parameter (if set
326     * to <tt>true</tt>) allows to ignore the maximum size and add the element to the queue.
327     *
328     * @return <tt>true</tt> if the element has been inserted, <tt>false</tt> if the element was not inserted (the queue
329     *         has reached its maximum size).
330     *
331     * @throws NullPointerException if the specified element is null
332     */
333    boolean offer(QueueElement<E> queueElement, boolean ignoreSize) {
334        if (queueElement == null) {
335            throw new NullPointerException("queueElement is NULL");
336        }
337        if (queueElement.getPriority() < 0 || queueElement.getPriority() >= priorities) {
338            throw new IllegalArgumentException("priority out of range: " + queueElement);
339        }
340        if (queueElement.inQueue) {
341            throw new IllegalStateException("queueElement already in a queue: " + queueElement);
342        }
343        if (!ignoreSize && currentSize != null && currentSize.get() >= maxSize) {
344            return false;
345        }
346        boolean accepted = queues[queueElement.getPriority()].offer(queueElement);
347        debug("offer([{0}]), to P[{1}] delay[{2}ms] accepted[{3}]", queueElement.getElement().toString(),
348              queueElement.getPriority(), queueElement.getDelay(TimeUnit.MILLISECONDS), accepted);
349        if (accepted) {
350            if (currentSize != null) {
351                currentSize.incrementAndGet();
352            }
353            queueElement.inQueue = true;
354        }
355        return accepted;
356    }
357
358    /**
359     * Insert the specified element into the queue.
360     * <p/>
361     * The element is added with minimun priority and no delay.
362     *
363     * @param queueElement the element to add.
364     *
365     * @return <tt>true</tt> if the element has been inserted, <tt>false</tt> if the element was not inserted (the queue
366     *         has reached its maximum size).
367     *
368     * @throws NullPointerException if the specified element is null
369     */
370    @Override
371    public boolean offer(QueueElement<E> queueElement) {
372        return offer(queueElement, false);
373    }
374
375    /**
376     * Retrieve and remove the head of this queue, or return <tt>null</tt> if this queue has no elements with an expired
377     * delay.
378     * <p/>
379     * The retrieved element is the oldest one from the highest priority sub-queue.
380     * <p/>
381     * Invocations to this method run the anti-starvation (once every interval check).
382     *
383     * @return the head of this queue, or <tt>null</tt> if this queue has no elements with an expired delay.
384     */
385    @Override
386    public QueueElement<E> poll() {
387        lock.lock();
388        try {
389            antiStarvation();
390            QueueElement<E> e = null;
391            int i = priorities;
392            for (; e == null && i > 0; i--) {
393                e = queues[i - 1].poll();
394            }
395            if (e != null) {
396                if (currentSize != null) {
397                    currentSize.decrementAndGet();
398                }
399                e.inQueue = false;
400                debug("poll(): [{0}], from P[{1}]", e.getElement().toString(), i);
401            }
402            return e;
403        }
404        finally {
405            lock.unlock();
406        }
407    }
408
409    /**
410     * Retrieve, but does not remove, the head of this queue, or returns <tt>null</tt> if this queue is empty.  Unlike
411     * <tt>poll</tt>, if no expired elements are available in the queue, this method returns the element that will
412     * expire next, if one exists.
413     *
414     * @return the head of this queue, or <tt>null</tt> if this queue is empty.
415     */
416    @Override
417    public QueueElement<E> peek() {
418        lock.lock();
419        try {
420            antiStarvation();
421            QueueElement<E> e = null;
422
423            QueueElement<E> [] seeks = new QueueElement[priorities];
424            boolean foundElement = false;
425            for (int i = priorities - 1; i > -1; i--) {
426                e = queues[i].peek();
427                debug("peek(): considering [{0}] from P[{1}]", e, i);
428                seeks[priorities - i - 1] = e;
429                foundElement |= e != null;
430            }
431            if (foundElement) {
432                e = null;
433                for (int i = 0; e == null && i < priorities; i++) {
434                    if (seeks[i] != null && seeks[i].getDelay(TimeUnit.MILLISECONDS) > 0) {
435                        debug("peek, ignoring [{0}]", seeks[i]);
436                    }
437                    else {
438                        e = seeks[i];
439                    }
440                }
441                if (e != null) {
442                    debug("peek(): choosing [{0}]", e);
443                }
444                if (e == null) {
445                    int first;
446                    for (first = 0; e == null && first < priorities; first++) {
447                        e = seeks[first];
448                    }
449                    if (e != null) {
450                        debug("peek(): initial choosing [{0}]", e);
451                    }
452                    for (int i = first; i < priorities; i++) {
453                        QueueElement<E> ee = seeks[i];
454                        if (ee != null && ee.getDelay(TimeUnit.MILLISECONDS) < e.getDelay(TimeUnit.MILLISECONDS)) {
455                            debug("peek(): choosing [{0}] over [{1}]", ee, e);
456                            e = ee;
457                        }
458                    }
459                }
460            }
461            if (e != null) {
462                debug("peek(): [{0}], from P[{1}]", e.getElement().toString(), e.getPriority());
463            }
464            else {
465                debug("peek(): NULL");
466            }
467            return e;
468        }
469        finally {
470            lock.unlock();
471        }
472    }
473
474    /**
475     * Run the anti-starvation check every {@link #ANTI_STARVATION_INTERVAL} milliseconds.
476     * <p/>
477     * It promotes elements beyond max wait time to the next higher priority sub-queue.
478     */
479    protected void antiStarvation() {
480        long now = System.currentTimeMillis();
481        if (now - lastAntiStarvationCheck > ANTI_STARVATION_INTERVAL) {
482            for (int i = 0; i < queues.length - 1; i++) {
483                antiStarvation(queues[i], queues[i + 1], "from P[" + i + "] to P[" + (i + 1) + "]");
484            }
485            StringBuilder sb = new StringBuilder();
486            for (int i = 0; i < queues.length; i++) {
487                sb.append("P[").append(i).append("]=").append(queues[i].size()).append(" ");
488            }
489            debug("sub-queue sizes: {0}", sb.toString());
490            lastAntiStarvationCheck = System.currentTimeMillis();
491        }
492    }
493
494    /**
495     * Promote elements beyond max wait time from a lower priority sub-queue to a higher priority sub-queue.
496     *
497     * @param lowerQ lower priority sub-queue.
498     * @param higherQ higher priority sub-queue.
499     * @param msg sub-queues msg (from-to) for debugging purposes.
500     */
501    private void antiStarvation(DelayQueue<QueueElement<E>> lowerQ, DelayQueue<QueueElement<E>> higherQ, String msg) {
502        int moved = 0;
503        QueueElement<E> e = lowerQ.poll();
504        while (e != null && e.getDelay(TimeUnit.MILLISECONDS) < -maxWait) {
505            e.setDelay(0, TimeUnit.MILLISECONDS);
506            if (!higherQ.offer(e)) {
507                throw new IllegalStateException("Could not move element to higher sub-queue, element rejected");
508            }
509            e.priority++;
510            e = lowerQ.poll();
511            moved++;
512        }
513        if (e != null) {
514            if (!lowerQ.offer(e)) {
515                throw new IllegalStateException("Could not reinsert element to current sub-queue, element rejected");
516            }
517        }
518        debug("anti-starvation, moved {0} element(s) {1}", moved, msg);
519    }
520
521    /**
522     * Method for debugging purposes. This implementation is a <tt>NOP</tt>.
523     * <p/>
524     * This method should be overriden for logging purposes.
525     * <p/>
526     * Message templates used by this class are in JDK's <tt>MessageFormat</tt> syntax.
527     *
528     * @param msgTemplate message template.
529     * @param msgArgs arguments for the message template.
530     */
531    protected void debug(String msgTemplate, Object... msgArgs) {
532    }
533
534    /**
535     * Insert the specified element into this queue, waiting if necessary
536     * for space to become available.
537     * <p/>
538     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
539     *
540     * @param e the element to add
541     * @throws InterruptedException if interrupted while waiting
542     * @throws ClassCastException if the class of the specified element
543     *         prevents it from being added to this queue
544     * @throws NullPointerException if the specified element is null
545     * @throws IllegalArgumentException if some property of the specified
546     *         element prevents it from being added to this queue
547     */
548    @Override
549    public void put(QueueElement<E> e) throws InterruptedException {
550        while (!offer(e, true)) {
551            Thread.sleep(10);
552        }
553    }
554
555    /**
556     * Insert the specified element into this queue, waiting up to the
557     * specified wait time if necessary for space to become available.
558     * <p/>
559     * IMPORTANT: This implementation forces the addition of the element to the queue regardless
560     * of the queue current size. The timeout value is ignored as the element is added immediately.
561     * <p/>
562     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
563     *
564     * @param e the element to add
565     * @param timeout how long to wait before giving up, in units of
566     *        <tt>unit</tt>
567     * @param unit a <tt>TimeUnit</tt> determining how to interpret the
568     *        <tt>timeout</tt> parameter
569     * @return <tt>true</tt> if successful, or <tt>false</tt> if
570     *         the specified waiting time elapses before space is available
571     * @throws InterruptedException if interrupted while waiting
572     * @throws ClassCastException if the class of the specified element
573     *         prevents it from being added to this queue
574     * @throws NullPointerException if the specified element is null
575     * @throws IllegalArgumentException if some property of the specified
576     *         element prevents it from being added to this queue
577     */
578    @Override
579    public boolean offer(QueueElement<E> e, long timeout, TimeUnit unit) throws InterruptedException {
580        return offer(e, true);
581    }
582
583    /**
584     * Retrieve and removes the head of this queue, waiting if necessary
585     * until an element becomes available.
586     * <p/>
587     * IMPORTANT: This implementation has a delay of up to 10ms (when the queue is empty) to detect a new element
588     * is available. It is doing a 10ms sleep.
589     * <p/>
590     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
591     *
592     * @return the head of this queue
593     * @throws InterruptedException if interrupted while waiting
594     */
595    @Override
596    public QueueElement<E> take() throws InterruptedException {
597        QueueElement<E> e = poll();
598        while (e == null) {
599            Thread.sleep(10);
600            e = poll();
601        }
602        return e;
603    }
604
605    /**
606     * Retrieve and removes the head of this queue, waiting up to the
607     * specified wait time if necessary for an element to become available.
608     * <p/>
609     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
610     *
611     * @param timeout how long to wait before giving up, in units of
612     *        <tt>unit</tt>
613     * @param unit a <tt>TimeUnit</tt> determining how to interpret the
614     *        <tt>timeout</tt> parameter
615     * @return the head of this queue, or <tt>null</tt> if the
616     *         specified waiting time elapses before an element is available
617     * @throws InterruptedException if interrupted while waiting
618     */
619    @Override
620    public QueueElement<E> poll(long timeout, TimeUnit unit) throws InterruptedException {
621        QueueElement<E> e = poll();
622        long time = System.currentTimeMillis() + unit.toMillis(timeout);
623        while (e == null && time > System.currentTimeMillis()) {
624            Thread.sleep(10);
625            e = poll();
626        }
627        return poll();
628    }
629
630    /**
631     * Return the number of additional elements that this queue can ideally
632     * (in the absence of memory or resource constraints) accept without
633     * blocking, or <tt>Integer.MAX_VALUE</tt> if there is no intrinsic
634     * limit.
635     *
636     * <p>Note that you <em>cannot</em> always tell if an attempt to insert
637     * an element will succeed by inspecting <tt>remainingCapacity</tt>
638     * because it may be the case that another thread is about to
639     * insert or remove an element.
640     * <p/>
641     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
642     *
643     * @return the remaining capacity
644     */
645    @Override
646    public int remainingCapacity() {
647        return (maxSize == -1) ? -1 : maxSize - size();
648    }
649
650    /**
651     * Remove all available elements from this queue and adds them
652     * to the given collection.  This operation may be more
653     * efficient than repeatedly polling this queue.  A failure
654     * encountered while attempting to add elements to
655     * collection <tt>c</tt> may result in elements being in neither,
656     * either or both collections when the associated exception is
657     * thrown.  Attempt to drain a queue to itself result in
658     * <tt>IllegalArgumentException</tt>. Further, the behavior of
659     * this operation is undefined if the specified collection is
660     * modified while the operation is in progress.
661     * <p/>
662     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
663     *
664     * @param c the collection to transfer elements into
665     * @return the number of elements transferred
666     * @throws UnsupportedOperationException if addition of elements
667     *         is not supported by the specified collection
668     * @throws ClassCastException if the class of an element of this queue
669     *         prevents it from being added to the specified collection
670     * @throws NullPointerException if the specified collection is null
671     * @throws IllegalArgumentException if the specified collection is this
672     *         queue, or some property of an element of this queue prevents
673     *         it from being added to the specified collection
674     */
675    @Override
676    public int drainTo(Collection<? super QueueElement<E>> c) {
677        int count = 0;
678        for (DelayQueue<QueueElement<E>> q : queues) {
679            count += q.drainTo(c);
680        }
681        return count;
682    }
683
684    /**
685     * Remove at most the given number of available elements from
686     * this queue and adds them to the given collection.  A failure
687     * encountered while attempting to add elements to
688     * collection <tt>c</tt> may result in elements being in neither,
689     * either or both collections when the associated exception is
690     * thrown.  Attempt to drain a queue to itself result in
691     * <tt>IllegalArgumentException</tt>. Further, the behavior of
692     * this operation is undefined if the specified collection is
693     * modified while the operation is in progress.
694     * <p/>
695     * NOTE: This method is to fulfill the <tt>BlockingQueue<tt/> interface. Not implemented in the most optimal way.
696     *
697     * @param c the collection to transfer elements into
698     * @param maxElements the maximum number of elements to transfer
699     * @return the number of elements transferred
700     * @throws UnsupportedOperationException if addition of elements
701     *         is not supported by the specified collection
702     * @throws ClassCastException if the class of an element of this queue
703     *         prevents it from being added to the specified collection
704     * @throws NullPointerException if the specified collection is null
705     * @throws IllegalArgumentException if the specified collection is this
706     *         queue, or some property of an element of this queue prevents
707     *         it from being added to the specified collection
708     */
709    @Override
710    public int drainTo(Collection<? super QueueElement<E>> c, int maxElements) {
711        int left = maxElements;
712        int count = 0;
713        for (DelayQueue<QueueElement<E>> q : queues) {
714            int drained = q.drainTo(c, left);
715            count += drained;
716            left -= drained;
717        }
718        return count;
719    }
720
721    /**
722     * Removes all of the elements from this queue. The queue will be empty after this call returns.
723     */
724    @Override
725    public void clear() {
726        for (DelayQueue<QueueElement<E>> q : queues) {
727            q.clear();
728        }
729    }
730}