Class SparkIEJoinOperator<Type0 extends java.lang.Comparable<Type0>,Type1 extends java.lang.Comparable<Type1>,Input extends Copyable<Input>>
- java.lang.Object
-
- org.apache.wayang.core.plan.wayangplan.OperatorBase
-
- org.apache.wayang.core.plan.wayangplan.BinaryToUnaryOperator<Input,Input,Tuple2<Input,Input>>
-
- org.apache.wayang.iejoin.operators.IEJoinOperator<Type0,Type1,Input>
-
- org.apache.wayang.iejoin.operators.SparkIEJoinOperator<Type0,Type1,Input>
-
- All Implemented Interfaces:
java.io.Serializable,ActualOperator,ElementaryOperator,ExecutionOperator,Operator,SparkExecutionOperator
public class SparkIEJoinOperator<Type0 extends java.lang.Comparable<Type0>,Type1 extends java.lang.Comparable<Type1>,Input extends Copyable<Input>> extends IEJoinOperator<Type0,Type1,Input> implements SparkExecutionOperator
Spark implementation of theIEJoinOperator.- 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.iejoin.operators.IEJoinOperator
cond0, cond1, equalReverse, get0Pivot, get0Ref, get1Pivot, get1Ref, list1ASC, list1ASCSec, list2ASC, list2ASCSec
-
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 SparkIEJoinOperator(DataSetType<Input> inputType0, DataSetType<Input> inputType1, TransformationDescriptor<Input,Type0> get0Pivot, TransformationDescriptor<Input,Type0> get1Pivot, IEJoinMasterOperator.JoinCondition cond0, TransformationDescriptor<Input,Type1> get0Ref, TransformationDescriptor<Input,Type1> get1Ref, IEJoinMasterOperator.JoinCondition cond1)Creates a new instance.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description booleancontainsAction()Tell whether this instances is a Spark action.protected ExecutionOperatorcreateCopy()Tuple<java.util.Collection<ExecutionLineageNode>,java.util.Collection<ChannelInstance>>evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, SparkExecutor sparkExecutor, OptimizationContext.OperatorContext operatorContext)Evaluates this operator.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.iejoin.operators.IEJoinOperator
assignSortOrders, getCond0, getCond1, getGet0Pivot, getGet0Ref, getGet1Pivot, getGet1Ref
-
Methods inherited from class org.apache.wayang.core.plan.wayangplan.BinaryToUnaryOperator
getInputType0, getInputType1, 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, createLoadProfileEstimator, createOutputChannelInstances, getLimitBaseKey, getLoadProfileEstimatorConfigurationKey, getLoadProfileEstimatorConfigurationKeys, getOriginal, getOutputChannelDescriptor, isFiltered
-
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
-
Methods inherited from interface org.apache.wayang.spark.operators.SparkExecutionOperator
getPlatform, name, name
-
-
-
-
Constructor Detail
-
SparkIEJoinOperator
public SparkIEJoinOperator(DataSetType<Input> inputType0, DataSetType<Input> inputType1, TransformationDescriptor<Input,Type0> get0Pivot, TransformationDescriptor<Input,Type0> get1Pivot, IEJoinMasterOperator.JoinCondition cond0, TransformationDescriptor<Input,Type1> get0Ref, TransformationDescriptor<Input,Type1> get1Ref, IEJoinMasterOperator.JoinCondition cond1)
Creates a new instance.
-
-
Method Detail
-
evaluate
public Tuple<java.util.Collection<ExecutionLineageNode>,java.util.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
protected ExecutionOperator createCopy()
- Overrides:
createCopyin classOperatorBase
-
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)
-
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
-
-