Package org.jppf.client.utils
Class AbstractJPPFJobStream
java.lang.Object
org.jppf.client.event.JobListenerAdapter
org.jppf.client.utils.AbstractJPPFJobStream
- All Implemented Interfaces:
AutoCloseable,Iterable<JPPFJob>,EventListener,Iterator<JPPFJob>,JobListener
public abstract class AbstractJPPFJobStream
extends JobListenerAdapter
implements Iterable<JPPFJob>, Iterator<JPPFJob>, AutoCloseable
Instances of this class provide a stream of JPPF jobs.
A common usage pattern is as follows:
// concurrency level
int concurrency = 4;
try (JPPFClient client = new JPPFClient();
AbstractJPPFJobStream jobStream = new MyJobStreamImplementation(concurrency)) {
jobStream.forEach(job -> client.submitJob(job));
jobsStream.awaitEndOfStream();
}-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionbooleanWait until this job stream has finished procesing all of its jobs.abstract voidclose()Close this stream and release the underlying resources it uses.protected abstract JPPFJobCreate the next job in the stream, along with its tasks.intGet the number of completed jobs.intGet the number of submitted jobs.intGet the number of submitted task.abstract booleanhasNext()Determine whether there is at least one more job in the stream.booleanDetermine whether any job is still being executed.iterator()voidThis implementation ofJobListener.jobEnded(JobEvent)decreases the counter of running jobs, notifies all threads waiting innext()and finally processes the results asynchronously.next()Get the next job in the stream.protected abstract voidprocessResults(JPPFJob job) Callback invoked when a job is complete.voidremove()This operation is not supported and results in anUnsupportedOperationExceptionbeing thrown.Methods inherited from class org.jppf.client.event.JobListenerAdapter
jobDispatched, jobReturned, jobStartedMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface java.lang.Iterable
forEach, spliteratorMethods inherited from interface java.util.Iterator
forEachRemaining
-
Constructor Details
-
AbstractJPPFJobStream
public AbstractJPPFJobStream(int concurrencyLimit) Initialize this job provider.- Parameters:
concurrencyLimit- the maximum number of jobs submitted concurrently.
-
-
Method Details
-
hasNext
public abstract boolean hasNext()Determine whether there is at least one more job in the stream. -
next
Get the next job in the stream.- Specified by:
nextin interfaceIterator<JPPFJob>- Returns:
- a newly created
JPPFJobobject. - Throws:
NoSuchElementException- if this stream has no more job to provide.
-
createNextJob
Create the next job in the stream, along with its tasks. This method must be overriden in subclasses. It does not need to update the internal state of this job stream, however it should manage the underlying resources it uses, such as files or database connections.This method is called each time
next()is invoked, including implicitely in enhancedforloops.- Returns:
- the created job.
-
remove
This operation is not supported and results in anUnsupportedOperationExceptionbeing thrown.- Specified by:
removein interfaceIterator<JPPFJob>- Throws:
UnsupportedOperationException- every time this method is called.
-
jobEnded
This implementation ofJobListener.jobEnded(JobEvent)decreases the counter of running jobs, notifies all threads waiting innext()and finally processes the results asynchronously.- Specified by:
jobEndedin interfaceJobListener- Overrides:
jobEndedin classJobListenerAdapter- Parameters:
event- encaspulates the source of the event.
-
processResults
Callback invoked when a job is complete.- Parameters:
job- the job whose results to process.
-
iterator
-
close
Close this stream and release the underlying resources it uses.- Specified by:
closein interfaceAutoCloseable- Throws:
Exception- if any error occurs.
-
hasPendingJob
public boolean hasPendingJob()Determine whether any job is still being executed.- Returns:
trueif at least one job was submitted and has not yet completed,falseotherwise.
-
getJobCount
public int getJobCount()Get the number of submitted jobs.- Returns:
- the count of submitted jobs.
-
getExecutedJobCount
public int getExecutedJobCount()Get the number of completed jobs.- Returns:
- the count of completed jobs.
-
getTaskCount
public int getTaskCount()Get the number of submitted task.- Returns:
- the count of submitted tasks.
-
awaitEndOfStream
public boolean awaitEndOfStream()Wait until this job stream has finished procesing all of its jobs.- Returns:
trueif the end of stream has been reached, false if the current thread was interrupted before the end of stream occurred.
-