Class FlinkGroupByOperator<InputType,KeyType>
- java.lang.Object
-
- org.apache.wayang.core.plan.wayangplan.OperatorBase
-
- org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator<Input,java.lang.Iterable<Input>>
-
- org.apache.wayang.basic.operators.GroupByOperator<InputType,KeyType>
-
- org.apache.wayang.flink.operators.FlinkGroupByOperator<InputType,KeyType>
-
- All Implemented Interfaces:
java.io.Serializable,ActualOperator,ElementaryOperator,ExecutionOperator,Operator,FlinkExecutionOperator
public class FlinkGroupByOperator<InputType,KeyType> extends GroupByOperator<InputType,KeyType> implements FlinkExecutionOperator
Flink implementation of theGroupByOperator.- 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.GroupByOperator
keyDescriptor
-
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 FlinkGroupByOperator(GroupByOperator<InputType,KeyType> that)Copies an instance (exclusive of broadcasts).FlinkGroupByOperator(TransformationDescriptor<InputType,KeyType> keydescriptor, DataSetType<InputType> inputType, DataSetType<KeyType> keyType)Creates a new instance.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description booleancontainsAction()Tell whether this instances is a Flink action.protected ExecutionOperatorcreateCopy()java.util.Optional<LoadProfileEstimator>createLoadProfileEstimator(Configuration configuration)Developers ofExecutionOperators can provide a defaultLoadProfileEstimatorvia this method.Tuple<java.util.Collection<ExecutionLineageNode>,java.util.Collection<ChannelInstance>>evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext)java.lang.StringgetLoadProfileEstimatorConfigurationKey()java.util.List<ChannelDescriptor>getSupportedInputChannels(int index)java.util.List<ChannelDescriptor>getSupportedOutputChannels(int index)Display the supportedChannels for a certainOutputSlot.-
Methods inherited from class org.apache.wayang.basic.operators.GroupByOperator
getKeyDescriptor
-
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
-
FlinkGroupByOperator
public FlinkGroupByOperator(TransformationDescriptor<InputType,KeyType> keydescriptor, DataSetType<InputType> inputType, DataSetType<KeyType> keyType)
Creates a new instance.
-
FlinkGroupByOperator
public FlinkGroupByOperator(GroupByOperator<InputType,KeyType> 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:
evaluatein interfaceFlinkExecutionOperator
-
containsAction
public boolean containsAction()
Description copied from interface:FlinkExecutionOperatorTell whether this instances is a Flink action. This is important to keep track on when Flink is actually initialized.- Specified by:
containsActionin interfaceFlinkExecutionOperator- Returns:
- whether this instance issues Flink actions
-
createCopy
protected ExecutionOperator createCopy()
- Overrides:
createCopyin classOperatorBase
-
getLoadProfileEstimatorConfigurationKey
public java.lang.String getLoadProfileEstimatorConfigurationKey()
- Specified by:
getLoadProfileEstimatorConfigurationKeyin interfaceExecutionOperator
-
createLoadProfileEstimator
public java.util.Optional<LoadProfileEstimator> createLoadProfileEstimator(Configuration configuration)
Description copied from interface:ExecutionOperatorDevelopers ofExecutionOperators can provide a defaultLoadProfileEstimatorvia this method.- Specified by:
createLoadProfileEstimatorin interfaceExecutionOperator- Parameters:
configuration- in which theLoadProfileshould be estimated.- Returns:
- an
Optionalthat might contain theLoadProfileEstimator(butOptional.empty()by default)
-
getSupportedInputChannels
public java.util.List<ChannelDescriptor> getSupportedInputChannels(int index)
Description copied from interface:ExecutionOperator- Specified by:
getSupportedInputChannelsin interfaceExecutionOperator- Parameters:
index- the index of theInputSlot- Returns:
- an
ListofChannels'Classes, ordered by their preference of use
-
getSupportedOutputChannels
public java.util.List<ChannelDescriptor> getSupportedOutputChannels(int index)
Description copied from interface:ExecutionOperatorDisplay the supportedChannels for a certainOutputSlot.- Specified by:
getSupportedOutputChannelsin interfaceExecutionOperator- Parameters:
index- the index of theOutputSlot- Returns:
- an
ListofChannels'Classes, ordered by their preference of use - See Also:
ExecutionOperator.getOutputChannelDescriptor(int)
-
-