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}