Class FlinkGlobalReduceOperator<Type>
- java.lang.Object
-
- org.apache.wayang.core.plan.wayangplan.OperatorBase
-
- org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator<Type,Type>
-
- org.apache.wayang.basic.operators.GlobalReduceOperator<Type>
-
- org.apache.wayang.flink.operators.FlinkGlobalReduceOperator<Type>
-
- All Implemented Interfaces:
java.io.Serializable
,ActualOperator
,ElementaryOperator
,ExecutionOperator
,Operator
,FlinkExecutionOperator
public class FlinkGlobalReduceOperator<Type> extends GlobalReduceOperator<Type> implements FlinkExecutionOperator
Flink implementation of theGlobalReduceOperator
.- 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 inherited from class org.apache.wayang.basic.operators.GlobalReduceOperator
reduceDescriptor
-
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 FlinkGlobalReduceOperator(GlobalReduceOperator<Type> that)
Copies an instance (exclusive of broadcasts).FlinkGlobalReduceOperator(DataSetType<Type> type, ReduceDescriptor<Type> reduceDescriptor)
Creates a new instance.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description boolean
containsAction()
Tell whether this instances is a Flink action.protected ExecutionOperator
createCopy()
java.util.Optional<LoadProfileEstimator>
createLoadProfileEstimator(Configuration configuration)
Developers ofExecutionOperator
s can provide a defaultLoadProfileEstimator
via this method.Tuple<java.util.Collection<ExecutionLineageNode>,java.util.Collection<ChannelInstance>>
evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext)
java.lang.String
getLoadProfileEstimatorConfigurationKey()
java.util.List<ChannelDescriptor>
getSupportedInputChannels(int index)
java.util.List<ChannelDescriptor>
getSupportedOutputChannels(int index)
Display the supportedChannel
s for a certainOutputSlot
.-
Methods inherited from class org.apache.wayang.basic.operators.GlobalReduceOperator
createCardinalityEstimator, getReduceDescriptor, getType
-
Methods inherited from class org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator
getInput, getInputType, getOutput, getOutputType
-
Methods inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
accept, addBroadcastInput, addTargetPlatform, at, collectMappedInputSlots, collectMappedOutputSlots, copy, 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
createCardinalityEstimator, getCardinalityEstimator, isAuxiliary, setAuxiliary, setCardinalityEstimator
-
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.ExecutionOperator
copy, createOutputChannelInstances, getLimitBaseKey, getLoadProfileEstimatorConfigurationKeys, getOriginal, getOutputChannelDescriptor, isFiltered
-
Methods inherited from interface org.apache.wayang.flink.operators.FlinkExecutionOperator
getBroadCastFunction, getPlatform
-
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
-
-
-
-
Constructor Detail
-
FlinkGlobalReduceOperator
public FlinkGlobalReduceOperator(DataSetType<Type> type, ReduceDescriptor<Type> reduceDescriptor)
Creates a new instance.- Parameters:
type
- type of the reduce elements (i.e., type ofUnaryToUnaryOperator.getInput()
andUnaryToUnaryOperator.getOutput()
)reduceDescriptor
- describes the reduction to be performed on the elements
-
FlinkGlobalReduceOperator
public FlinkGlobalReduceOperator(GlobalReduceOperator<Type> that)
Copies an instance (exclusive of broadcasts).- Parameters:
that
- that should be copied
-
-
Method Detail
-
evaluate
public Tuple<java.util.Collection<ExecutionLineageNode>,java.util.Collection<ChannelInstance>> evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext)
- Specified by:
evaluate
in interfaceFlinkExecutionOperator
-
createCopy
protected ExecutionOperator createCopy()
- Overrides:
createCopy
in classOperatorBase
-
getLoadProfileEstimatorConfigurationKey
public java.lang.String getLoadProfileEstimatorConfigurationKey()
- Specified by:
getLoadProfileEstimatorConfigurationKey
in interfaceExecutionOperator
-
createLoadProfileEstimator
public java.util.Optional<LoadProfileEstimator> createLoadProfileEstimator(Configuration configuration)
Description copied from interface:ExecutionOperator
Developers ofExecutionOperator
s can provide a defaultLoadProfileEstimator
via this method.- Specified by:
createLoadProfileEstimator
in interfaceExecutionOperator
- Parameters:
configuration
- in which theLoadProfile
should be estimated.- Returns:
- an
Optional
that might contain theLoadProfileEstimator
(butOptional.empty()
by default)
-
getSupportedInputChannels
public java.util.List<ChannelDescriptor> getSupportedInputChannels(int index)
Description copied from interface:ExecutionOperator
- Specified by:
getSupportedInputChannels
in interfaceExecutionOperator
- Parameters:
index
- the index of theInputSlot
- Returns:
- an
List
ofChannel
s'Class
es, ordered by their preference of use
-
getSupportedOutputChannels
public java.util.List<ChannelDescriptor> getSupportedOutputChannels(int index)
Description copied from interface:ExecutionOperator
Display the supportedChannel
s for a certainOutputSlot
.- Specified by:
getSupportedOutputChannels
in interfaceExecutionOperator
- Parameters:
index
- the index of theOutputSlot
- Returns:
- an
List
ofChannel
s'Class
es, ordered by their preference of use - See Also:
ExecutionOperator.getOutputChannelDescriptor(int)
-
containsAction
public boolean containsAction()
Description copied from interface:FlinkExecutionOperator
Tell whether this instances is a Flink action. This is important to keep track on when Flink is actually initialized.- Specified by:
containsAction
in interfaceFlinkExecutionOperator
- Returns:
- whether this instance issues Flink actions
-
-