public class StatefulParDoP<OutputT>
extends java.lang.Object
Processor implementation for Beam's stateful ParDo primitive.| Modifier and Type | Class and Description | 
|---|---|
| static class  | StatefulParDoP.Supplier<OutputT>Jet  Processorsupplier that will provide instances ofStatefulParDoP. | 
| Modifier and Type | Method and Description | 
|---|---|
| void | close() | 
| boolean | complete() | 
| boolean | completeEdge(int ordinal) | 
| protected void | finishRunnerBundle(org.apache.beam.runners.core.DoFnRunner<InputT,OutputT> runner) | 
| protected org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> | getDoFnRunner(PipelineOptions pipelineOptions,
             DoFn<KV<?,?>,OutputT> doFn,
             org.apache.beam.runners.core.SideInputReader sideInputReader,
             org.apache.beam.runners.jet.processors.AbstractParDoP.JetOutputManager outputManager,
             TupleTag<OutputT> mainOutputTag,
             java.util.List<TupleTag<?>> additionalOutputTags,
             Coder<KV<?,?>> inputValueCoder,
             java.util.Map<TupleTag<?>,Coder<?>> outputValueCoders,
             WindowingStrategy<?,?> windowingStrategy,
             DoFnSchemaInformation doFnSchemaInformation,
             java.util.Map<java.lang.String,PCollectionView<?>> sideInputMapping) | 
| void | init(com.hazelcast.jet.core.Outbox outbox,
    com.hazelcast.jet.core.Processor.Context context) | 
| boolean | isCooperative() | 
| void | process(int ordinal,
       com.hazelcast.jet.core.Inbox inbox) | 
| protected void | processElementWithRunner(org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> runner,
                        org.apache.beam.sdk.util.WindowedValue<KV<?,?>> windowedValue) | 
| protected void | startRunnerBundle(org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> runner) | 
| boolean | tryProcess() | 
| boolean | tryProcessWatermark(com.hazelcast.jet.core.Watermark watermark) | 
protected org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> getDoFnRunner(PipelineOptions pipelineOptions, DoFn<KV<?,?>,OutputT> doFn, org.apache.beam.runners.core.SideInputReader sideInputReader, org.apache.beam.runners.jet.processors.AbstractParDoP.JetOutputManager outputManager, TupleTag<OutputT> mainOutputTag, java.util.List<TupleTag<?>> additionalOutputTags, Coder<KV<?,?>> inputValueCoder, java.util.Map<TupleTag<?>,Coder<?>> outputValueCoders, WindowingStrategy<?,?> windowingStrategy, DoFnSchemaInformation doFnSchemaInformation, java.util.Map<java.lang.String,PCollectionView<?>> sideInputMapping)
protected void startRunnerBundle(org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> runner)
protected void processElementWithRunner(org.apache.beam.runners.core.DoFnRunner<KV<?,?>,OutputT> runner, org.apache.beam.sdk.util.WindowedValue<KV<?,?>> windowedValue)
public boolean tryProcessWatermark(@Nonnull com.hazelcast.jet.core.Watermark watermark)
tryProcessWatermark in interface com.hazelcast.jet.core.Processorpublic boolean complete()
complete in interface com.hazelcast.jet.core.Processorpublic void init(@Nonnull com.hazelcast.jet.core.Outbox outbox, @Nonnull com.hazelcast.jet.core.Processor.Context context)
init in interface com.hazelcast.jet.core.Processorpublic boolean isCooperative()
isCooperative in interface com.hazelcast.jet.core.Processorpublic void close()
close in interface com.hazelcast.jet.core.Processorpublic void process(int ordinal,
                    @Nonnull
                    com.hazelcast.jet.core.Inbox inbox)
process in interface com.hazelcast.jet.core.Processorprotected void finishRunnerBundle(org.apache.beam.runners.core.DoFnRunner<InputT,OutputT> runner)
public boolean tryProcess()
tryProcess in interface com.hazelcast.jet.core.Processorpublic boolean completeEdge(int ordinal)
completeEdge in interface com.hazelcast.jet.core.Processor