Package org.jppf.execute.async
Class AbstractAsyncExecutionManager
java.lang.Object
org.jppf.execute.async.AbstractAsyncExecutionManager
- All Implemented Interfaces:
AsyncExecutionManager
Instances of this class manage the execution of JPPF tasks by a node.
- Author:
- Laurent Cohen, Martin JANDA, Paul Woodward
-
Field Summary
FieldsModifier and TypeFieldDescriptionprotected final AtomicBooleanDetermines whether the number of threads or their priority has changed.protected final CollectionMap<String,Long> Mapping of job uuids to the ids of the bundles currently processed for this job.protected final Map<String,JobProcessingEntry> Mapping of jobUuid + bunldeId to the correspondingJobProcessingEntryobjects.protected final List<ExecutionManagerListener>List of listeners to this execution manager.protected final CollectionMap<String,Long> Mapping of job uuids to the ids of the pending bundles for each job.protected final Map<String,JobPendingEntry> Map of the bundles being read from the server and not yet submitted to this execution manager.protected AtomicReference<JPPFReconnectionNotification>Set if the node must reconnect to the driver.protected final TaskExecutionDispatcherDispatches tasks notifications to registered listeners.protected final ThreadManagerThe thread manager that is used for execution.protected final JPPFScheduleHandlerTimer managing the tasks timeout. -
Constructor Summary
ConstructorsConstructorDescriptionAbstractAsyncExecutionManager(TypedProperties config, JPPFProperty<Integer> nbThreadsProperty) Initialize this execution manager with the specified node. -
Method Summary
Modifier and TypeMethodDescriptionvoidRegister a listener with this execution manager.voidaddPendingJobEntry(TaskBundle bundle) Called when a task bundle s being received by the node, and before it is submitted to this execution manager.voidcancelAllTasks(boolean callOnCancel, boolean requeue) Cancel all executing or pending tasks.voidCancel all executing or pending tasks for the specified job.booleanDetermines whether the configuration has changed and resets the flag if it has.protected abstract voidcleanup(JobProcessingEntry jobEntry) Cleanup method invoked when all tasks for the current bundle have completed.voidexecute(BundleWithTasks bundleWithTasks) Execute the specified tasks of the specified tasks bundle.protected voidfireJobFinished(TaskBundle bundle, List<Task<?>> tasks, Throwable t) Called when the execution of a task bundle has finished.Get the executor used by this execution manager.intgetNbBundles(String jobUuid) Get the number of bundles with the specified job uuid that are pending or being processed.Get the object which dispatches tasks notifications to registered listeners.Get the thread manager for this node.intGet the size of the node's thread pool.intGet the priority assigned to the execution threads.voidRemove a listener from the registered listeners.voidsetThreadPoolSize(int size) Set the size of the node's thread pool.protected abstract JobProcessingEntrysetup(BundleWithTasks bundleWithTasks) Prepare this execution manager for executing the tasks of a bundle.voidshutdown()Shutdown this execution manager.voidtaskEnded(NodeTaskWrapper taskWrapper) Notification that a task has finished executing.voidTrigger the configuration changed flag.voidupdateThreadsPriority(int newPriority) Update the priority of all execution threads.
-
Field Details
-
timeoutHandler
Timer managing the tasks timeout. -
taskNotificationDispatcher
Dispatches tasks notifications to registered listeners. -
configChanged
Determines whether the number of threads or their priority has changed. -
reconnectionNotification
Set if the node must reconnect to the driver. -
threadManager
The thread manager that is used for execution. -
jobEntries
Mapping of jobUuid + bunldeId to the correspondingJobProcessingEntryobjects. -
jobBundleIds
Mapping of job uuids to the ids of the bundles currently processed for this job. -
listeners
List of listeners to this execution manager. -
pendingEntries
Map of the bundles being read from the server and not yet submitted to this execution manager. -
pendingBundleIds
Mapping of job uuids to the ids of the pending bundles for each job.
-
-
Constructor Details
-
AbstractAsyncExecutionManager
public AbstractAsyncExecutionManager(TypedProperties config, JPPFProperty<Integer> nbThreadsProperty) Initialize this execution manager with the specified node.- Parameters:
config- the configuration to get the thread manager properties from.nbThreadsProperty- the name of the property which configures the number of threads.
-
-
Method Details
-
execute
Description copied from interface:AsyncExecutionManagerExecute the specified tasks of the specified tasks bundle.- Specified by:
executein interfaceAsyncExecutionManager- Parameters:
bundleWithTasks- the bundle to execute with its associated tasks.- Throws:
Exception- if the execution failed.
-
cancelAllTasks
public void cancelAllTasks(boolean callOnCancel, boolean requeue) Description copied from interface:AsyncExecutionManagerCancel all executing or pending tasks.- Specified by:
cancelAllTasksin interfaceAsyncExecutionManager- Parameters:
callOnCancel- determines whether the onCancel() callback method of each task should be invoked.requeue- true if the job should be requeued on the server side, false otherwise.
-
cancelJob
Description copied from interface:AsyncExecutionManagerCancel all executing or pending tasks for the specified job.- Specified by:
cancelJobin interfaceAsyncExecutionManager- Parameters:
jobUuid- the uuid of the job to cancel.callOnCancel- determines whether the onCancel() callback method of each task should be invoked.requeue- true if the job should be requeued on the server side, false otherwise.
-
shutdown
public void shutdown()Description copied from interface:AsyncExecutionManagerShutdown this execution manager.- Specified by:
shutdownin interfaceAsyncExecutionManager
-
setup
Prepare this execution manager for executing the tasks of a bundle.- Parameters:
bundleWithTasks- the bundle and associated tasks.- Returns:
- an instance of
JobProcessingEntry.
-
cleanup
Cleanup method invoked when all tasks for the current bundle have completed.- Parameters:
jobEntry- encapsulates information about the job.
-
taskEnded
Description copied from interface:AsyncExecutionManagerNotification that a task has finished executing.- Specified by:
taskEndedin interfaceAsyncExecutionManager- Parameters:
taskWrapper- the task that finished its execution.
-
getExecutor
Description copied from interface:AsyncExecutionManagerGet the executor used by this execution manager.- Specified by:
getExecutorin interfaceAsyncExecutionManager- Returns:
- an
ExecutorServiceinstance.
-
checkConfigChanged
public boolean checkConfigChanged()Description copied from interface:AsyncExecutionManagerDetermines whether the configuration has changed and resets the flag if it has.- Specified by:
checkConfigChangedin interfaceAsyncExecutionManager- Returns:
- true if the config was changed, false otherwise.
-
triggerConfigChanged
public void triggerConfigChanged()Description copied from interface:AsyncExecutionManagerTrigger the configuration changed flag.- Specified by:
triggerConfigChangedin interfaceAsyncExecutionManager
-
setThreadPoolSize
public void setThreadPoolSize(int size) Description copied from interface:AsyncExecutionManagerSet the size of the node's thread pool.- Specified by:
setThreadPoolSizein interfaceAsyncExecutionManager- Parameters:
size- the size as an int.
-
getThreadPoolSize
public int getThreadPoolSize()Description copied from interface:AsyncExecutionManagerGet the size of the node's thread pool.- Specified by:
getThreadPoolSizein interfaceAsyncExecutionManager- Returns:
- the size as an int.
-
getThreadsPriority
public int getThreadsPriority()Description copied from interface:AsyncExecutionManagerGet the priority assigned to the execution threads.- Specified by:
getThreadsPriorityin interfaceAsyncExecutionManager- Returns:
- the priority as an int value.
-
updateThreadsPriority
public void updateThreadsPriority(int newPriority) Description copied from interface:AsyncExecutionManagerUpdate the priority of all execution threads.- Specified by:
updateThreadsPriorityin interfaceAsyncExecutionManager- Parameters:
newPriority- the new priority to set.
-
getThreadManager
Description copied from interface:AsyncExecutionManagerGet the thread manager for this node.- Specified by:
getThreadManagerin interfaceAsyncExecutionManager- Returns:
- a
ThreadManagerinstance.
-
getTaskNotificationDispatcher
Description copied from interface:AsyncExecutionManagerGet the object which dispatches tasks notifications to registered listeners.- Specified by:
getTaskNotificationDispatcherin interfaceAsyncExecutionManager- Returns:
- a
TaskExecutionDispatcherinstance.
-
getNbBundles
Description copied from interface:AsyncExecutionManagerGet the number of bundles with the specified job uuid that are pending or being processed.- Specified by:
getNbBundlesin interfaceAsyncExecutionManager- Parameters:
jobUuid- the uuid of the job for which to count the bundles.- Returns:
- number of bundles.
-
addExecutionManagerListener
Description copied from interface:AsyncExecutionManagerRegister a listener with this execution manager.- Specified by:
addExecutionManagerListenerin interfaceAsyncExecutionManager- Parameters:
listener- the listener to register.
-
removeExecutionManagerListener
Description copied from interface:AsyncExecutionManagerRemove a listener from the registered listeners.- Specified by:
removeExecutionManagerListenerin interfaceAsyncExecutionManager- Parameters:
listener- the listener to remove.
-
fireJobFinished
Called when the execution of a task bundle has finished.- Parameters:
bundle- the TaskBundle which holds information on the job.tasks- the tasks that were executed.t- aThrowablethat prevented or interrupted the job processing.
-
addPendingJobEntry
Description copied from interface:AsyncExecutionManagerCalled when a task bundle s being received by the node, and before it is submitted to this execution manager.- Specified by:
addPendingJobEntryin interfaceAsyncExecutionManager- Parameters:
bundle- the bundle to process.
-