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
Processor supplier that will provide instances of StatefulParDoP. |
| 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