Class SparkGlobalReduceOperator<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.spark.operators.SparkGlobalReduceOperator<Type>
- All Implemented Interfaces:
Serializable,ActualOperator,ElementaryOperator,ExecutionOperator,Operator,SparkExecutionOperator
public class SparkGlobalReduceOperator<Type>
extends GlobalReduceOperator<Type>
implements SparkExecutionOperator
Spark implementation of the
GlobalReduceOperator.- See Also:
-
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
reduceDescriptorFields inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
inputSlots, outputSlots, STANDARD_OPERATOR_ARGSFields inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
FIRST_EPOCH -
Constructor Summary
ConstructorsConstructorDescriptionCopies an instance (exclusive of broadcasts).SparkGlobalReduceOperator(DataSetType<Type> type, ReduceDescriptor<Type> reduceDescriptor) Creates a new instance. -
Method Summary
Modifier and TypeMethodDescriptionbooleanTell whether this instances is a Spark action.protected ExecutionOperatorcreateLoadProfileEstimator(Configuration configuration) Developers ofExecutionOperators can provide a defaultLoadProfileEstimatorvia this method.evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, SparkExecutor sparkExecutor, OptimizationContext.OperatorContext operatorContext) Evaluates this operator.getSupportedInputChannels(int index) getSupportedOutputChannels(int index) Display the supportedChannels for a certainOutputSlot.Methods inherited from class org.apache.wayang.basic.operators.GlobalReduceOperator
createCardinalityEstimator, getReduceDescriptor, getTypeMethods inherited from class org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator
getInput, getInputType, getOutput, getOutputTypeMethods 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, toStringMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, waitMethods inherited from interface org.apache.wayang.core.plan.wayangplan.ActualOperator
acceptMethods inherited from interface org.apache.wayang.core.plan.wayangplan.ElementaryOperator
createCardinalityEstimator, getCardinalityEstimator, isAuxiliary, setAuxiliary, setCardinalityEstimatorMethods inherited from interface org.apache.wayang.core.plan.wayangplan.ExecutionOperator
copy, createOutputChannelInstances, getLimitBaseKey, getLoadProfileEstimatorConfigurationKeys, getOriginal, getOutputChannelDescriptor, isFilteredMethods 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, isConversion, isElementary, isExecutionOperator, isFeedbackInput, isFeedforwardOutput, isLoopHead, isLoopSubplan, isOwnerOf, isReading, isSink, isSource, isSubplan, isSupportingBroadcastInputs, isUnconnected, propagateInputCardinality, propagateOutputCardinality, propagateOutputCardinality, replaceWith, setContainer, setEpoch, setInput, setName, setOutputMethods inherited from interface org.apache.wayang.spark.operators.SparkExecutionOperator
getPlatform, name, name
-
Constructor Details
-
SparkGlobalReduceOperator
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
-
SparkGlobalReduceOperator
Copies an instance (exclusive of broadcasts).- Parameters:
that- that should be copied
-
-
Method Details
-
evaluate
public Tuple<Collection<ExecutionLineageNode>,Collection<ChannelInstance>> evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, SparkExecutor sparkExecutor, OptimizationContext.OperatorContext operatorContext) Description copied from interface:SparkExecutionOperatorEvaluates this operator. Takes a set ofChannelInstances according to the operator inputs and manipulates a set ofChannelInstances according to the operator outputs -- unless the operator is a sink, then it triggers execution.In addition, this method should give feedback of what this instance was doing by wiring the
LazyExecutionLineageNodes of input and ouputChannelInstances and providing aCollectionof executedExecutionLineageNodes.- Specified by:
evaluatein interfaceSparkExecutionOperator- Parameters:
inputs-ChannelInstances that satisfy the inputs of this operatoroutputs-ChannelInstances that accept the outputs of this operatorsparkExecutor-SparkExecutorthat executes this instanceoperatorContext- optimization information for this instance- Returns:
Collections of what has been executed and produced
-
createCopy
- Overrides:
createCopyin classOperatorBase
-
getLoadProfileEstimatorConfigurationKey
- Specified by:
getLoadProfileEstimatorConfigurationKeyin interfaceExecutionOperator
-
createLoadProfileEstimator
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
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
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:
-
containsAction
public boolean containsAction()Description copied from interface:SparkExecutionOperatorTell whether this instances is a Spark action. This is important to keep track on when Spark is actually initialized.- Specified by:
containsActionin interfaceSparkExecutionOperator- Returns:
- whether this instance issues Spark actions
-