diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupQueue.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupQueue.java index d4c3473492..e60a975f12 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupQueue.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupQueue.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2011 the original author or authors. + * Copyright 2002-2012 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -36,14 +36,16 @@ import org.springframework.util.Assert; * the face of transaction rollback (assuming the store is transactional) and also to ensure messages are not lost if * the process dies (assuming the store is durable). To use the queue across process re-starts, the same group id * must be provided, so it needs to be unique but identifiable with a single logical instance of the queue. - * + * * @author Dave Syer * @author Oleg Zhurakousky + * @author Gunnar Hillert + * * @since 2.0 - * + * */ public class MessageGroupQueue extends AbstractQueue> implements BlockingQueue> { - + private final Log logger = LogFactory.getLog(getClass()); private static final int DEFAULT_CAPACITY = Integer.MAX_VALUE; @@ -53,14 +55,14 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc private final Object groupId; private final int capacity; - + //This one could be a global semaphore private final Lock storeLock; - - private final Condition messageStoreNotFull; - - private final Condition messageStoreNotEmpty; - + + private final Condition messageStoreNotFull; + + private final Condition messageStoreNotEmpty; + public MessageGroupQueue(MessageGroupStore messageGroupStore, Object groupId) { this(messageGroupStore, groupId, DEFAULT_CAPACITY, new ReentrantLock(true)); } @@ -68,11 +70,11 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc public MessageGroupQueue(MessageGroupStore messageGroupStore, Object groupId, int capacity) { this(messageGroupStore, groupId, capacity, new ReentrantLock(true)); } - + public MessageGroupQueue(MessageGroupStore messageGroupStore, Object groupId, Lock storeLock) { this(messageGroupStore, groupId, DEFAULT_CAPACITY, storeLock); } - + public MessageGroupQueue(MessageGroupStore messageGroupStore, Object groupId, int capacity, Lock storeLock) { Assert.isTrue(capacity > 0, "'capacity' must be greater than 0"); Assert.notNull(storeLock, "'storeLock' must not be null"); @@ -104,30 +106,30 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc if (!messages.isEmpty()) { message = messages.iterator().next(); } - } + } finally { storeLock.unlock(); } - } + } catch (InterruptedException e) { Thread.currentThread().interrupt(); } return message; } - + public Message poll(long timeout, TimeUnit unit) throws InterruptedException { Message message = null; long timeoutInNanos = unit.toNanos(timeout); final Lock storeLock = this.storeLock; storeLock.lockInterruptibly(); - - try { + + try { while (this.size() == 0 && timeoutInNanos > 0){ - timeoutInNanos = this.messageStoreNotEmpty.awaitNanos(timeoutInNanos); + timeoutInNanos = this.messageStoreNotEmpty.awaitNanos(timeoutInNanos); } message = this.doPoll(); - - } + + } finally { storeLock.unlock(); } @@ -141,17 +143,17 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc storeLock.lockInterruptibly(); try { message = this.doPoll(); - } + } finally { storeLock.unlock(); } - } + } catch (InterruptedException e) { Thread.currentThread().interrupt(); } return message; } - + public int drainTo(Collection> c) { return this.drainTo(c, Integer.MAX_VALUE); } @@ -163,18 +165,18 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc final Lock storeLock = this.storeLock; try { storeLock.lockInterruptibly(); - try { + try { Message message = this.messageGroupStore.pollMessageFromGroup(groupId); for (int i = 0; i < maxElements && message != null; i++) { - list.add(message); + list.add(message); message = this.messageGroupStore.pollMessageFromGroup(groupId); } this.messageStoreNotFull.signal(); - } + } finally { storeLock.unlock(); - } - } + } + } catch (InterruptedException e) { logger.warn("Queue may not have drained completely since this operation was interrupted", e); Thread.currentThread().interrupt(); @@ -182,19 +184,19 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc collection.addAll(list); return collection.size() - originalSize; } - + public boolean offer(Message message) { boolean offered = true; - final Lock storeLock = this.storeLock; + final Lock storeLock = this.storeLock; try { storeLock.lockInterruptibly(); - try { + try { offered = this.doOffer(message); - } + } finally { storeLock.unlock(); - } - } + } + } catch (InterruptedException e) { Thread.currentThread().interrupt(); } @@ -204,40 +206,45 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc public boolean offer(Message message, long timeout, TimeUnit unit) throws InterruptedException { long timeoutInNanos = unit.toNanos(timeout); boolean offered = false; - + final Lock storeLock = this.storeLock; storeLock.lockInterruptibly(); try { - while (this.size() == capacity && timeoutInNanos > 0){ - timeoutInNanos = this.messageStoreNotFull.awaitNanos(timeoutInNanos); + if (capacity != Integer.MAX_VALUE) { + while (this.size() == capacity && timeoutInNanos > 0){ + timeoutInNanos = this.messageStoreNotFull.awaitNanos(timeoutInNanos); + } } - if (timeoutInNanos > 0){ offered = this.doOffer(message); } - } + } finally { storeLock.unlock(); } - return offered; + return offered; } public void put(Message message) throws InterruptedException { final Lock storeLock = this.storeLock; storeLock.lockInterruptibly(); try { - while (this.size() == capacity){ - this.messageStoreNotFull.await(); + if (capacity != Integer.MAX_VALUE) { + while (this.size() == capacity){ + this.messageStoreNotFull.await(); + } } - this.doOffer(message); - } + } finally { storeLock.unlock(); } } public int remainingCapacity() { + if (capacity == Integer.MAX_VALUE) { + return Integer.MAX_VALUE; + } return capacity - this.size(); } @@ -245,14 +252,14 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc Message message = null; final Lock storeLock = this.storeLock; storeLock.lockInterruptibly(); - - try { + + try { while (this.size() == 0){ - this.messageStoreNotEmpty.await(); + this.messageStoreNotEmpty.await(); } message = this.doPoll(); - - } + + } finally { storeLock.unlock(); } @@ -272,14 +279,14 @@ public class MessageGroupQueue extends AbstractQueue> implements Bloc this.messageStoreNotFull.signal(); return message; } - + /** * It is assumed that the 'storeLock' is being held by the caller, otherwise * IllegalMonitorStateException may be thrown */ private boolean doOffer(Message message){ boolean offered = false; - if (this.size() < capacity){ + if (capacity == Integer.MAX_VALUE || this.size() < capacity){ messageGroupStore.addMessageToGroup(groupId, message); offered = true; this.messageStoreNotEmpty.signal();