Class LocalCallbackSink<T>
- java.lang.Object
-
- org.apache.wayang.core.plan.wayangplan.OperatorBase
-
- org.apache.wayang.core.plan.wayangplan.UnarySink<T>
-
- org.apache.wayang.basic.operators.LocalCallbackSink<T>
-
- All Implemented Interfaces:
java.io.Serializable
,ActualOperator
,ElementaryOperator
,Operator
- Direct Known Subclasses:
FlinkLocalCallbackSink
,JavaLocalCallbackSink
,SparkLocalCallbackSink
public class LocalCallbackSink<T> extends UnarySink<T>
This sink executes a callback on each received data unit into a JavaCollection
.- See Also:
- Serialized Form
-
-
Nested Class Summary
-
Nested classes/interfaces inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
OperatorBase.GsonSerializer
-
-
Field Summary
Fields Modifier and Type Field Description protected java.util.function.Consumer<T>
callback
protected FunctionDescriptor.SerializableConsumer<T>
callbackDescriptor
protected java.util.Collection<T>
collector
-
Fields inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
inputSlots, outputSlots, STANDARD_OPERATOR_ARGS
-
Fields inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
FIRST_EPOCH
-
-
Constructor Summary
Constructors Constructor Description LocalCallbackSink(java.lang.Class<T> typeClass)
Convnience constructor, defaults to StdoutSinkLocalCallbackSink(java.util.function.Consumer<T> callback, java.lang.Class<T> typeClass)
Creates a new instance.LocalCallbackSink(java.util.function.Consumer<T> callback, DataSetType<T> type)
Creates a new instance.LocalCallbackSink(LocalCallbackSink<T> that)
Copies an instance (exclusive of broadcasts).LocalCallbackSink(FunctionDescriptor.SerializableConsumer<T> consumerDescriptor, java.lang.Class<T> typeClass)
Creates a new instance.LocalCallbackSink(FunctionDescriptor.SerializableConsumer<T> consumerDescriptor, DataSetType<T> type)
Creates a new instance.
-
Method Summary
All Methods Static Methods Instance Methods Concrete Methods Modifier and Type Method Description java.util.Optional<CardinalityEstimator>
createCardinalityEstimator(int outputIndex, Configuration configuration)
static <T> LocalCallbackSink<T>
createCollectingSink(java.util.Collection<T> collector, java.lang.Class<T> typeClass)
static <T> LocalCallbackSink<T>
createCollectingSink(java.util.Collection<T> collector, DataSetType<T> type)
static <T> LocalCallbackSink<T>
createStdoutSink(java.lang.Class<T> typeClass)
static <T> LocalCallbackSink<T>
createStdoutSink(DataSetType<T> type)
java.util.function.Consumer<T>
getCallback()
FunctionDescriptor.SerializableConsumer<T>
getCallbackDescriptor()
LocalCallbackSink<T>
setCollector(java.util.Collection<T> collector)
-
Methods inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
accept, addBroadcastInput, addTargetPlatform, at, collectMappedInputSlots, collectMappedOutputSlots, copy, createCopy, getAllInputs, getAllOutputs, getCardinalityEstimator, getContainer, getEpoch, getName, getOriginal, getSimpleClassName, getTargetPlatforms, isAuxiliary, isSupportingBroadcastInputs, propagateInputCardinality, propagateOutputCardinality, setAuxiliary, setCardinalityEstimator, setContainer, setEpoch, setName, toString
-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, wait
-
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.ActualOperator
accept
-
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.ElementaryOperator
getCardinalityEstimator, isAuxiliary, setAuxiliary, setCardinalityEstimator
-
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
addBroadcastInput, addTargetPlatform, broadcastTo, broadcastTo, collectMappedInputSlots, collectMappedOutputSlots, connectTo, connectTo, getAllInputs, getAllOutputs, getCardinalityPusher, getContainer, getEffectiveOccupant, getEffectiveOccupant, getEpoch, getEstimationContextProperties, getForwards, getInnermostLoop, getInput, getInput, getLoopStack, getName, getNumBroadcastInputs, getNumInputs, getNumOutputs, getNumRegularInputs, getOuterInputSlot, getOutermostInputSlot, getOutermostOutputSlots, getOutput, getOutput, getParent, getTargetPlatforms, isAlternative, isElementary, isExecutionOperator, isFeedbackInput, isFeedforwardOutput, isLoopHead, isLoopSubplan, isOwnerOf, isReading, isSink, isSource, isSubplan, isSupportingBroadcastInputs, isUnconnected, propagateInputCardinality, propagateOutputCardinality, propagateOutputCardinality, setContainer, setEpoch, setInput, setName, setOutput
-
-
-
-
Field Detail
-
callback
protected final java.util.function.Consumer<T> callback
-
callbackDescriptor
protected final FunctionDescriptor.SerializableConsumer<T> callbackDescriptor
-
collector
protected java.util.Collection<T> collector
-
-
Constructor Detail
-
LocalCallbackSink
public LocalCallbackSink(java.util.function.Consumer<T> callback, DataSetType<T> type)
Creates a new instance.- Parameters:
callback
- callback that is executed locally for each incoming data unittype
- type of the incoming elements
-
LocalCallbackSink
public LocalCallbackSink(LocalCallbackSink<T> that)
Copies an instance (exclusive of broadcasts).- Parameters:
that
- that should be copied
-
LocalCallbackSink
public LocalCallbackSink(java.util.function.Consumer<T> callback, java.lang.Class<T> typeClass)
Creates a new instance.- Parameters:
callback
- callback that is executed locally for each incoming data unittypeClass
- type of the incoming elements
-
LocalCallbackSink
public LocalCallbackSink(FunctionDescriptor.SerializableConsumer<T> consumerDescriptor, java.lang.Class<T> typeClass)
Creates a new instance.
-
LocalCallbackSink
public LocalCallbackSink(FunctionDescriptor.SerializableConsumer<T> consumerDescriptor, DataSetType<T> type)
Creates a new instance.- Parameters:
type
- type of the dataunit elements
-
LocalCallbackSink
public LocalCallbackSink(java.lang.Class<T> typeClass)
Convnience constructor, defaults to StdoutSink
-
-
Method Detail
-
createCollectingSink
public static <T> LocalCallbackSink<T> createCollectingSink(java.util.Collection<T> collector, DataSetType<T> type)
-
createCollectingSink
public static <T> LocalCallbackSink<T> createCollectingSink(java.util.Collection<T> collector, java.lang.Class<T> typeClass)
-
createStdoutSink
public static <T> LocalCallbackSink<T> createStdoutSink(DataSetType<T> type)
-
createStdoutSink
public static <T> LocalCallbackSink<T> createStdoutSink(java.lang.Class<T> typeClass)
-
setCollector
public LocalCallbackSink<T> setCollector(java.util.Collection<T> collector)
-
getCallback
public java.util.function.Consumer<T> getCallback()
-
getCallbackDescriptor
public FunctionDescriptor.SerializableConsumer<T> getCallbackDescriptor()
-
createCardinalityEstimator
public java.util.Optional<CardinalityEstimator> createCardinalityEstimator(int outputIndex, Configuration configuration)
Description copied from interface:ElementaryOperator
- Parameters:
outputIndex
- index of theOutputSlot
for that theCardinalityEstimator
is requestedconfiguration
- if theCardinalityEstimator
depends on further ones, use this to obtain the latter- Returns:
- an
Optional
that might provide the requested instance
-
-