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 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.
      Specified by:
      hasNext in interface Iterator<JPPFJob>
      Returns:
      true if there is at least one job in the stream, false otherwise.
    • next

      public JPPFJob next() throws NoSuchElementException
      Get the next job in the stream.
      Specified by:
      next in interface Iterator<JPPFJob>
      Returns:
      a newly created JPPFJob object.
      Throws:
      NoSuchElementException - if this stream has no more job to provide.
    • createNextJob

      protected abstract JPPFJob 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 enhanced for loops.

      Returns:
      the created job.
    • remove

      public void remove() throws UnsupportedOperationException
      This operation is not supported and results in an UnsupportedOperationException being thrown.
      Specified by:
      remove in interface Iterator<JPPFJob>
      Throws:
      UnsupportedOperationException - every time this method is called.
    • jobEnded

      public void jobEnded(JobEvent event)
      This implementation of JobListener.jobEnded(JobEvent) decreases the counter of running jobs, notifies all threads waiting in next() and finally processes the results asynchronously.
      Specified by:
      jobEnded in interface JobListener
      Overrides:
      jobEnded in class JobListenerAdapter
      Parameters:
      event - encaspulates the source of the event.
    • processResults

      protected abstract void processResults(JPPFJob job)
      Callback invoked when a job is complete.
      Parameters:
      job - the job whose results to process.
    • iterator

      public Iterator<JPPFJob> iterator()
      Specified by:
      iterator in interface Iterable<JPPFJob>
    • close

      public abstract void close() throws Exception
      Close this stream and release the underlying resources it uses.
      Specified by:
      close in interface AutoCloseable
      Throws:
      Exception - if any error occurs.
    • hasPendingJob

      public boolean hasPendingJob()
      Determine whether any job is still being executed.
      Returns:
      true if at least one job was submitted and has not yet completed, false otherwise.
    • 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:
      true if the end of stream has been reached, false if the current thread was interrupted before the end of stream occurred.