Class JPPFExecutorService
- All Implemented Interfaces:
Executor,ExecutorService,EventListener,JobListener
ExecutorService wrapper around a JPPFClient.
This executor has two modes in which it functions:
1) Standard mode: in this mode each task or set of tasks submitted via one of the
invokeXXX() or submit() methods is sent immediately to the server in its own JPPF job.
2) Batch mode: the JPPFExecutorService can be configured to only send tasks to the server
when a number of tasks, submitted via one of the invokeXXX() or submit() methods,
has been reached, or when a timeout specified in milliseconds has expired, or a combination of both.
This facility is designed to optimize the task execution throughput, especially when many individual tasks are submitted
using one of the submit() methods. This way, the tasks are sent to the server as a single job,
instead of one job per task, and the execution will fully benefit from the parallel features of the JPPF server, including
scheduling, load-balancing and parallel I/O.
In batch mode, the following behavior is to be noted:
- If both size-based and time-based batching are used, tasks will be sent whenever one of the two thresholds is reached. Whenever this happens, both counters are reset. For instance, if the size-based threshold is reached, then the time-based counter will be reset as well, and the timeout counting will start from 0 again
- When a collection of tasks is submitted via one of the
invokeXXX()methods, they are guaranteed to be all sent together in the same JPPF job. This is the one exception to the batch size threshold. - If one of the threshold is changed while tasks are still pending execution, the behavior is unspecified
- Author:
- Laurent Cohen
- See Also:
-
Constructor Summary
ConstructorsConstructorDescriptionJPPFExecutorService(JPPFClient client) Initialize this executor service with the specified JPPF client.JPPFExecutorService(JPPFClient client, int batchSize, long batchTimeout) Initialize this executor service with the specified JPPF client, batch size and batch tiemout. -
Method Summary
Modifier and TypeMethodDescriptionbooleanawaitTermination(long timeout, TimeUnit unit) Blocks until all tasks have completed execution after a shutdown request, or the timeout occurs, or the current thread is interrupted, whichever happens first.voidExecutes the given command at some time in the future.intGet the minimum number of tasks that must be submitted before they are sent to the server.longGet the maximum time to wait before the next batch of tasks is to be sent for execution.Get the configuration for this executor service.invokeAll(Collection<? extends Callable<T>> tasks) Executes the given tasks, returning a list of Futures holding their status and results when all complete.invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) Executes the given tasks, returning a list of Futures holding their status and results when all complete or the timeout expires, whichever happens first.<T> TinvokeAny(Collection<? extends Callable<T>> tasks) Executes the given tasks, returning the result of one that has completed successfully (i.e., without throwing an exception), if any do.<T> TinvokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) Executes the given tasks, returning the result of one that has completed successfully (i.e., without throwing an exception), if any do before the given timeout elapses.booleanDetermine whether this executor has been shut down.booleanDetermine whether all tasks have completed following shut down.voidjobReturned(JobEvent event) Called when all results from a job have been received.Reset the configuration for this executor service to a blank state.setBatchSize(int batchSize) Set the minimum number of tasks that must be submitted before they are sent to the server.setBatchTimeout(long batchTimeout) Set the maximum time to wait before the next batch of tasks is to be sent for execution.voidshutdown()Initiates an orderly shutdown in which previously submitted tasks are executed, but no new tasks will be accepted.Attempts to stop all actively executing tasks, halts the processing of waiting tasks, and returns a list of the tasks that were awaiting execution.
This implementation simply waits for all submitted tasks to terminate, due to the complexity of stopping remote tasks.Future<?>Submits a Runnable task for execution and returns a Future representing that task.<T> Future<T>Submits a Runnable task for execution and returns a Future representing that task that will upon completion return the given result.<T> Future<T>Submit a value-returning task for execution and returns a Future representing the pending results of the task.Methods inherited from class org.jppf.client.event.JobListenerAdapter
jobDispatched, jobEnded, jobStarted
-
Constructor Details
-
JPPFExecutorService
Initialize this executor service with the specified JPPF client.- Parameters:
client- theJPPFClientto use for job submission.
-
JPPFExecutorService
Initialize this executor service with the specified JPPF client, batch size and batch tiemout.- Parameters:
client- theJPPFClientto use for job submission.batchSize- the minimum number of tasks that must be submitted before they are sent to the server.batchTimeout- the maximum time to wait before the next batch of tasks is to be sent for execution.
-
-
Method Details
-
invokeAll
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException Executes the given tasks, returning a list of Futures holding their status and results when all complete.- Specified by:
invokeAllin interfaceExecutorService- Type Parameters:
T- the type of results returned by the tasks.- Parameters:
tasks- the tasks to execute.- Returns:
- a list of Futures representing the tasks, in the same sequential order as produced by the iterator for the given task list, each of which has completed.
- Throws:
InterruptedException- if interrupted while waiting, in which case unfinished tasks are cancelled.NullPointerException- if tasks or any of its elements are null.RejectedExecutionException- if any task cannot be scheduled for execution.
-
invokeAll
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException Executes the given tasks, returning a list of Futures holding their status and results when all complete or the timeout expires, whichever happens first.- Specified by:
invokeAllin interfaceExecutorService- Type Parameters:
T- the type of results returned by the tasks.- Parameters:
tasks- the tasks to execute.timeout- the maximum time to wait.unit- the time unit of the timeout argument.- Returns:
- a list of Futures representing the tasks, in the same sequential order as produced by the iterator for the given task list, each of which has completed.
- Throws:
InterruptedException- if interrupted while waiting, in which case unfinished tasks are cancelled.NullPointerException- if tasks or any of its elements are null.RejectedExecutionException- if any task cannot be scheduled for execution.
-
invokeAny
public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException Executes the given tasks, returning the result of one that has completed successfully (i.e., without throwing an exception), if any do. Upon normal or exceptional return, tasks that have not completed are cancelled.- Specified by:
invokeAnyin interfaceExecutorService- Type Parameters:
T- the type of results returned by the tasks.- Parameters:
tasks- the tasks to execute.- Returns:
- the result returned by one of the tasks.
- Throws:
InterruptedException- if interrupted while waiting.NullPointerException- if tasks or any of its elements are null.IllegalArgumentException- if tasks empty.ExecutionException- if no task successfully completes.RejectedExecutionException- if tasks cannot be scheduled for execution.
-
invokeAny
public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException Executes the given tasks, returning the result of one that has completed successfully (i.e., without throwing an exception), if any do before the given timeout elapses. Upon normal or exceptional return, tasks that have not completed are cancelled.- Specified by:
invokeAnyin interfaceExecutorService- Type Parameters:
T- the type of results returned by the tasks.- Parameters:
tasks- the tasks to execute.timeout- the maximum time to wait.unit- the time unit of the timeout argument.- Returns:
- the result returned by one of the tasks.
- Throws:
InterruptedException- if interrupted while waiting.NullPointerException- if tasks or any of its elements are null.IllegalArgumentException- if tasks empty.ExecutionException- if no task successfully completes.RejectedExecutionException- if tasks cannot be scheduled for execution.TimeoutException- if the given timeout elapses before any task successfully completes.
-
submit
Submit a value-returning task for execution and returns a Future representing the pending results of the task.- Specified by:
submitin interfaceExecutorService- Type Parameters:
T- the type of result returned by the task.- Parameters:
task- the task to execute.- Returns:
- a Future representing pending completion of the task.
-
submit
Submits a Runnable task for execution and returns a Future representing that task.- Specified by:
submitin interfaceExecutorService- Parameters:
task- the task to execute.- See Also:
-
submit
Submits a Runnable task for execution and returns a Future representing that task that will upon completion return the given result.- Specified by:
submitin interfaceExecutorService- Type Parameters:
T- the type of result returned by the task.- Parameters:
task- the task to execute.result- the result to return .- Returns:
- a Future representing pending completion of the task, and whose get() method will return the given result upon completion.
-
execute
Executes the given command at some time in the future. The command may execute in a new thread, in a pooled thread, or in the calling thread, at the discretion of the Executor implementation. -
awaitTermination
Blocks until all tasks have completed execution after a shutdown request, or the timeout occurs, or the current thread is interrupted, whichever happens first.- Specified by:
awaitTerminationin interfaceExecutorService- Parameters:
timeout- the maximum time to wait.unit- the time unit of the timeout argument.- Returns:
- true if this executor terminated and false if the timeout elapsed before termination.
- Throws:
InterruptedException- if interrupted while waiting.
-
isShutdown
public boolean isShutdown()Determine whether this executor has been shut down.- Specified by:
isShutdownin interfaceExecutorService- Returns:
- true if this executor has been shut down, false otherwise.
-
isTerminated
public boolean isTerminated()Determine whether all tasks have completed following shut down. Note that isTerminated is never true unless either shutdown or shutdownNow was called first.- Specified by:
isTerminatedin interfaceExecutorService- Returns:
- true if all tasks have completed following shut down.
-
shutdown
public void shutdown()Initiates an orderly shutdown in which previously submitted tasks are executed, but no new tasks will be accepted.- Specified by:
shutdownin interfaceExecutorService
-
shutdownNow
Attempts to stop all actively executing tasks, halts the processing of waiting tasks, and returns a list of the tasks that were awaiting execution.
This implementation simply waits for all submitted tasks to terminate, due to the complexity of stopping remote tasks.- Specified by:
shutdownNowin interfaceExecutorService- Returns:
- a list of tasks that never commenced execution.
-
jobReturned
Called when all results from a job have been received.- Specified by:
jobReturnedin interfaceJobListener- Overrides:
jobReturnedin classJobListenerAdapter- Parameters:
event- the event object.
-
getBatchSize
public int getBatchSize()Get the minimum number of tasks that must be submitted before they are sent to the server.- Returns:
- the batch size as an int.
-
setBatchSize
Set the minimum number of tasks that must be submitted before they are sent to the server.- Parameters:
batchSize- the batch size as an int.- Returns:
- this executor service, for method chaining.
-
getBatchTimeout
public long getBatchTimeout()Get the maximum time to wait before the next batch of tasks is to be sent for execution.- Returns:
- the timeout as a long.
-
setBatchTimeout
Set the maximum time to wait before the next batch of tasks is to be sent for execution.- Parameters:
batchTimeout- the timeout as a long.- Returns:
- this executor service, for method chaining.
-
getConfiguration
Get the configuration for this executor service.- Returns:
- an
ExecutorServiceConfigurationinstance.
-
resetConfiguration
Reset the configuration for this executor service to a blank state.- Returns:
- an
ExecutorServiceConfigurationinstance.
-