Class FlinkDoWhileOperator<InputType,ConvergenceType>
java.lang.Object
org.apache.wayang.core.plan.wayangplan.OperatorBase
org.apache.wayang.basic.operators.DoWhileOperator<InputType,ConvergenceType>
org.apache.wayang.flink.operators.FlinkDoWhileOperator<InputType,ConvergenceType>
- All Implemented Interfaces:
Serializable
,ActualOperator
,ElementaryOperator
,ExecutionOperator
,LoopHeadOperator
,Operator
,FlinkExecutionOperator
public class FlinkDoWhileOperator<InputType,ConvergenceType>
extends DoWhileOperator<InputType,ConvergenceType>
implements FlinkExecutionOperator
Flink implementation of the
DoWhileOperator
.- See Also:
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
OperatorBase.GsonSerializer
Nested classes/interfaces inherited from interface org.apache.wayang.core.plan.wayangplan.LoopHeadOperator
LoopHeadOperator.State
-
Field Summary
Fields inherited from class org.apache.wayang.basic.operators.DoWhileOperator
CONVERGENCE_INPUT_INDEX, criterionDescriptor, FINAL_OUTPUT_INDEX, INITIAL_INPUT_INDEX, ITERATION_INPUT_INDEX, ITERATION_OUTPUT_INDEX
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
ConstructorsConstructorDescriptionCreates a new instance.FlinkDoWhileOperator
(DataSetType<InputType> inputType, DataSetType<ConvergenceType> convergenceType, FunctionDescriptor.SerializablePredicate<Collection<ConvergenceType>> criterionPredicate, Integer numExpectedIterations) Creates a new instance.FlinkDoWhileOperator
(DataSetType<InputType> inputType, DataSetType<ConvergenceType> convergenceType, PredicateDescriptor<Collection<ConvergenceType>> criterionDescriptor, Integer numExpectedIterations) -
Method Summary
Modifier and TypeMethodDescriptionboolean
Tell whether this instances is a Flink action.protected ExecutionOperator
createLoadProfileEstimator
(Configuration configuration) Developers ofExecutionOperator
s can provide a defaultLoadProfileEstimator
via this method.evaluate
(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext) getSupportedInputChannels
(int index) getSupportedOutputChannels
(int index) Display the supportedChannel
s for a certainOutputSlot
.Methods inherited from class org.apache.wayang.basic.operators.DoWhileOperator
beginIteration, createCardinalityEstimator, endIteration, getConditionInputSlots, getConditionOutputSlots, getConvergenceType, getCriterionDescriptor, getFinalLoopOutputs, getForwards, getInputType, getLoopBodyInputs, getLoopBodyOutputs, getLoopInitializationInputs, getNumExpectedIterations, getState, initialize, isReading, outputConnectTo, setNumExpectedIterations, setState
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.LoopHeadOperator
getCardinalityPusher, getFinalizationPusher, getInitializationPusher
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
addBroadcastInput, addTargetPlatform, broadcastTo, broadcastTo, collectMappedInputSlots, collectMappedOutputSlots, connectTo, connectTo, getAllInputs, getAllOutputs, getContainer, getEffectiveOccupant, getEffectiveOccupant, getEpoch, getEstimationContextProperties, getForwards, getInnermostLoop, getInput, getInput, getLoopStack, getName, getNumBroadcastInputs, getNumInputs, getNumOutputs, getNumRegularInputs, getOuterInputSlot, getOutermostInputSlot, getOutermostOutputSlots, getOutput, getOutput, getParent, getTargetPlatforms, isAlternative, isConversion, isElementary, isExecutionOperator, isFeedbackInput, isFeedforwardOutput, isLoopHead, isLoopSubplan, isOwnerOf, isReading, isSink, isSource, isSubplan, isSupportingBroadcastInputs, isUnconnected, propagateInputCardinality, propagateOutputCardinality, propagateOutputCardinality, replaceWith, setContainer, setEpoch, setInput, setName, setOutput
-
Constructor Details
-
FlinkDoWhileOperator
public FlinkDoWhileOperator(DataSetType<InputType> inputType, DataSetType<ConvergenceType> convergenceType, FunctionDescriptor.SerializablePredicate<Collection<ConvergenceType>> criterionPredicate, Integer numExpectedIterations) Creates a new instance. -
FlinkDoWhileOperator
public FlinkDoWhileOperator(DataSetType<InputType> inputType, DataSetType<ConvergenceType> convergenceType, PredicateDescriptor<Collection<ConvergenceType>> criterionDescriptor, Integer numExpectedIterations) -
FlinkDoWhileOperator
Creates a new instance.
-
-
Method Details
-
evaluate
public Tuple<Collection<ExecutionLineageNode>,Collection<ChannelInstance>> evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext) - Specified by:
evaluate
in interfaceFlinkExecutionOperator
-
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
-
getLoadProfileEstimatorConfigurationKey
- Specified by:
getLoadProfileEstimatorConfigurationKey
in interfaceExecutionOperator
-
createLoadProfileEstimator
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)
-
createCopy
- Overrides:
createCopy
in classOperatorBase
-
getSupportedInputChannels
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
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:
-