- Type Parameters:
T- type of emitted items
- All Implemented Interfaces:
public class ParallelStreamP<T> extends AbstractProcessor
Nested Class Summary
Nested classes/interfaces inherited from class com.hazelcast.jet.core.AbstractProcessor
Modifier and Type Method Description
()Called after all the inbound edges' streams are exhausted.
Processor.Context context)(Method that can be overridden to perform any necessary initialization for the processor.
Methods inherited from class com.hazelcast.jet.core.AbstractProcessor
emitFromTraverser, emitFromTraverser, emitFromTraverser, emitFromTraverserToSnapshot, flatMapper, flatMapper, flatMapper, getLogger, getOutbox, init, process, restoreFromSnapshot, restoreFromSnapshot, tryEmit, tryEmit, tryEmit, tryEmitToSnapshot, tryProcess, tryProcess0, tryProcess1, tryProcess2, tryProcess3, tryProcess4, tryProcessWatermark
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
ParallelStreamPpublic ParallelStreamP(long eventsPerSecondPerGenerator, EventTimePolicy<? super T> eventTimePolicy, List<? extends GeneratorFunction<T>> generators)Creates a processor that generates items using its assigned generator functions. This processor picks its generator functions according to its global processor index.
eventsPerSecondPerGenerator- the desired event rate for each generator
eventTimePolicy- parameters for watermark generation
generators- list of generator functions used in source
initDescription copied from class:
AbstractProcessorMethod that can be overridden to perform any necessary initialization for the processor. It is called exactly once and strictly before any of the processing methods (
Processor.complete()), but after the outbox and
loggerhave been initialized.
Subclasses are not required to call this superclass method, it does nothing.
completepublic boolean complete()Description copied from interface:
ProcessorCalled after all the inbound edges' streams are exhausted. If it returns
false, it will be invoked again until it returns
true. For example, a streaming source processor will return
falseforever. Unlike other methods which guarantee that no other method is called until they return
Processor.saveToSnapshot()can be called even though this method returned
After this method is called, no other processing methods are called on this processor, except for
Non-cooperative processors are required to return from this method from time to time to give the system a chance to check for snapshot requests and job cancellation. The time the processor spends in this method affects the latency of snapshots and job cancellations.
trueif the completing step is now done,
falseto call this method again