diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalState.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalState.java index 84a040f34..0b50d910d 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalState.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalState.java @@ -18,8 +18,20 @@ package org.springframework.batch.repeat.support; import java.util.Collection; +/** + * Internal interface for extensions of {@link RepeatTemplate}. + * + * @author Dave Syer + * + */ public interface RepeatInternalState { + /** + * Returns a mutable collection of exceptions that have occurred in the + * current repeat context. Clients are expected to mutate this collection. + * + * @return the collection of exceptions being accumulated + */ Collection getThrowables(); } diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalStateSupport.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalStateSupport.java index aefc38bdd..16c844e04 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalStateSupport.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatInternalStateSupport.java @@ -23,8 +23,8 @@ import java.util.Set; public class RepeatInternalStateSupport implements RepeatInternalState { // Accumulation of failed results. - private Set throwables = new HashSet(); - + private final Set throwables = new HashSet(); + /* (non-Javadoc) * @see org.springframework.batch.repeat.support.BatchInternalState#getThrowables() */ @@ -32,5 +32,4 @@ public class RepeatInternalStateSupport implements RepeatInternalState { return throwables; } - } diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java index 7a72ae195..594124c87 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/RepeatTemplate.java @@ -324,6 +324,8 @@ public class RepeatTemplate implements RepeatOperations { * * @param context the current {@link RepeatContext} * @return a {@link RepeatInternalState} instance. + * + * @see RepeatTemplate#waitForResults(RepeatInternalState) */ protected RepeatInternalState createInternalState(RepeatContext context) { return new RepeatInternalStateSupport(); @@ -331,7 +333,10 @@ public class RepeatTemplate implements RepeatOperations { /** * Get the next completed result, possibly executing several callbacks until - * one finally finishes. + * one finally finishes. Normally a subclass would have to override both + * this method and {@link #createInternalState(RepeatContext)} because the + * implementation of this method would rely on the details of the internal + * state. * * @param context current BatchContext. * @param callback the callback to execute. @@ -339,6 +344,7 @@ public class RepeatTemplate implements RepeatOperations { * @return a finished result. * * @see #isComplete(RepeatContext) + * @see #createInternalState(RepeatContext) */ protected ExitStatus getNextResult(RepeatContext context, RepeatCallback callback, RepeatInternalState state) throws Throwable { diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueue.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueue.java index c69cd203c..3ab47a630 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueue.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueue.java @@ -16,26 +16,66 @@ package org.springframework.batch.repeat.support; +import org.springframework.core.task.TaskExecutor; /** - * Abstraction for queue of {@link ResultHolder} objects. Acts as a BlockingQueue with the ability to count the number - * of items it expects to ever hold. When clients schedule an item to be added they call {@link #expect()}, and then - * when the result is collected the queue is notified that it no longer expects another. + * Abstraction for queue of {@link ResultHolder} objects. Acts as a + * BlockingQueue with the ability to count the number of items it expects to + * ever hold. When clients schedule an item to be added they call + * {@link #expect()}, and then collect the result later with {@link #take()}. + * Result providers in another thread call {@link #put(Object)} to notify the + * expecting client of a new result. * * @author Dave Syer * @author Ben Hale - * - * TODO: BlockingQueue with a CompletionService? */ -interface ResultQueue extends RepeatInternalState { +interface ResultQueue { + /** + * In a master-slave pattern, the master calls this method paired with + * {@link #take()} to manage the flow of items. Normally a task is submitted + * for processing in another thread, at which point the master uses this + * method to keep track of the number of expected results. It has the + * personality of an counter increment, rather than a work queue, which is + * usually managed elsewhere, e.g. by a {@link TaskExecutor}.

+ * Implementations may choose to block here, if they need to limit the + * number or rate of tasks being submitted. + * @throws InterruptedException if the call blocks and is then interrupted. + */ + void expect() throws InterruptedException; + + /** + * In a master-worker pattern, the workers call this method to deposit the + * result of a finished task on the queue for collection. + * + * @param result the result for later collection. + */ + void put(T result); + + /** + * Gets the next available result, blocking if there are none. + * + * @return a result previously deposited + * @throws InterruptedException if the operation is interrupted while + * waiting + */ + T take() throws InterruptedException; + + /** + * Used by master thread to verify that there are results available from + * {@link #take()} without possibly having to block and wait. + * + * @return true if there are no results available + */ boolean isEmpty(); - ResultHolder take(); - - void expect(); - - void put(ResultHolder holder); - + /** + * Check if any results are expected. Usually used by master thread to drain + * queue when it is finished. + * + * @return true if more results are expected, but possibly not yet + * available. + */ public boolean isExpecting(); + } diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueueFactory.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueueFactory.java deleted file mode 100644 index c6e59da49..000000000 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ResultQueueFactory.java +++ /dev/null @@ -1,32 +0,0 @@ -/* - * Copyright 2002-2007 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. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.batch.repeat.support; - -/** - * A factory for {@link ResultQueue} which simply creates one from - * java.util.concurrent components. - * - * @author Ben Hale - * @author Dave Syer - */ -class ResultQueueFactory { - - public RepeatInternalState getResultQueue(int throttleLimit) { - return new JdkConcurrentResultQueue(throttleLimit); - } - -} diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/TaskExecutorRepeatTemplate.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/TaskExecutorRepeatTemplate.java index 08964c058..15ef3f63b 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/TaskExecutorRepeatTemplate.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/TaskExecutorRepeatTemplate.java @@ -19,6 +19,7 @@ package org.springframework.batch.repeat.support; import org.springframework.batch.repeat.ExitStatus; import org.springframework.batch.repeat.RepeatCallback; import org.springframework.batch.repeat.RepeatContext; +import org.springframework.batch.repeat.RepeatException; import org.springframework.batch.repeat.RepeatOperations; import org.springframework.core.task.SyncTaskExecutor; import org.springframework.core.task.TaskExecutor; @@ -82,7 +83,7 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { ExecutingRunnable runnable = null; - ResultQueue queue = (ResultQueue) state; + ResultQueue queue = ((ResultQueueInternalState) state).getResultQueue(); do { @@ -137,7 +138,7 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { */ protected boolean waitForResults(RepeatInternalState state) { - ResultQueue queue = (ResultQueue) state; + ResultQueue queue = ((ResultQueueInternalState) state).getResultQueue(); boolean result = true; @@ -147,7 +148,14 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { * Careful that no runnables that are not going to finish ever get * onto the queue, else this may block forever. */ - ResultHolder future = (ResultHolder) queue.take(); + ResultHolder future; + try { + future = (ResultHolder) queue.take(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RepeatException("InterruptedException while waiting for result."); + } if (future.getError() != null) { state.getThrowables().add(future.getError()); @@ -167,7 +175,7 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { protected RepeatInternalState createInternalState(RepeatContext context) { // Queue of pending results: - return new ResultQueueFactory().getResultQueue(throttleLimit); + return new ResultQueueInternalState(throttleLimit); } /** @@ -181,13 +189,13 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { private RepeatContext context; - private ResultQueue queue; + private ResultQueue queue; private ExitStatus result; private Throwable error; - public ExecutingRunnable(RepeatCallback callback, RepeatContext context, ResultQueue queue) { + public ExecutingRunnable(RepeatCallback callback, RepeatContext context, ResultQueue queue) { super(); @@ -201,7 +209,13 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { * Tell the queue to expect a result. */ public void expect() { - queue.expect(); + try { + queue.expect(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RepeatException("InterruptedException waiting for to acquire lock on input."); + } } /** @@ -247,6 +261,30 @@ public class TaskExecutorRepeatTemplate extends RepeatTemplate { } + /** + * @author Dave Syer + * + */ + private static class ResultQueueInternalState extends RepeatInternalStateSupport { + + private final ResultQueue results; + + /** + * @param throttleLimit the throttle limit for the result queue + */ + public ResultQueueInternalState(int throttleLimit) { + super(); + this.results = new ThrottleLimitResultQueue(throttleLimit); + } + /* (non-Javadoc) + * @see org.springframework.batch.repeat.support.RepeatInternalState#getResultQueue() + */ + public ResultQueue getResultQueue() { + return results; + } + + } + /** * Public setter for the throttle limit. The throttle limit is the largest * number of concurrent tasks that can be executing at one time - if a new diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/JdkConcurrentResultQueue.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ThrottleLimitResultQueue.java similarity index 53% rename from spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/JdkConcurrentResultQueue.java rename to spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ThrottleLimitResultQueue.java index 47cd52d97..7309c67be 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/JdkConcurrentResultQueue.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/repeat/support/ThrottleLimitResultQueue.java @@ -20,17 +20,16 @@ import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.Semaphore; -import org.springframework.batch.repeat.RepeatException; - /** - * An implementation of the {@link ResultQueue} that uses the Java 5 Concurrent Utilities. + * An implementation of the {@link ResultQueue} that throttles the number of + * expected results, limiting it ti a maximum at any given time. * - * @author Ben Hale + * @author Dave Syer */ -class JdkConcurrentResultQueue extends RepeatInternalStateSupport implements ResultQueue { +class ThrottleLimitResultQueue implements ResultQueue { // Accumulation of result objects as they finish. - private final BlockingQueue results; + private final BlockingQueue results; // Accumulation of dummy objects flagging expected results in the future. private final Semaphore waits; @@ -39,31 +38,22 @@ class JdkConcurrentResultQueue extends RepeatInternalStateSupport implements Res private volatile int count = 0; - JdkConcurrentResultQueue(int throttleLimit) { - results = new LinkedBlockingQueue(); + /** + * @param throttleLimit the maximum number of results that can be expected + * at any given time. + */ + public ThrottleLimitResultQueue(int throttleLimit) { + results = new LinkedBlockingQueue(); waits = new Semaphore(throttleLimit); } - protected void addResult(ResultHolder resultHolder) { - results.add(resultHolder); - } - - protected void aquireWait() throws InterruptedException { - waits.acquire(); - } - - protected void releaseWait() { - waits.release(); - } - - protected ResultHolder takeResult() throws InterruptedException { - return (ResultHolder) results.take(); - } - public boolean isEmpty() { return results.isEmpty(); } + /* (non-Javadoc) + * @see org.springframework.batch.repeat.support.ResultQueue#isExpecting() + */ public boolean isExpecting() { synchronized (lock) { // Base the decision about whether we expect more results on a @@ -72,39 +62,34 @@ class JdkConcurrentResultQueue extends RepeatInternalStateSupport implements Res } } - public void expect() { - try { - synchronized (lock) { - aquireWait(); - count++; - } - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new RepeatException("InterruptedException waiting for to acquire lock on input."); + /** + * Tell the queue to expect one more result. Blocks until a new result is + * available if already expecting too many (as determined by the throttle + * limit). + * + * @see ResultQueue#expect() + */ + public void expect() throws InterruptedException { + synchronized (lock) { + waits.acquire(); + count++; } } - public void put(ResultHolder holder) { + public void put(T holder) { // There should be no need to block here, or to use offer() - addResult(holder); + results.add(holder); // Take from the waits queue now to allow another result to // accumulate. But don't decrement the counter. - releaseWait(); + waits.release(); } - public ResultHolder take() { - ResultHolder value; - try { - synchronized (lock) { - value = takeResult(); - // Decrement the counter only when the result is collected. - count--; - } - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new RepeatException("InterruptedException while waiting for result."); + public T take() throws InterruptedException { + T value; + synchronized (lock) { + value = results.take(); + // Decrement the counter only when the result is collected. + count--; } return value; } diff --git a/spring-batch-samples/foo.txt b/spring-batch-samples/foo.txt new file mode 100644 index 000000000..e965047ad --- /dev/null +++ b/spring-batch-samples/foo.txt @@ -0,0 +1 @@ +Hello