Class FlinkCoGroupOperator<InputType0,InputType1,TypeKey>
java.lang.Object
org.apache.wayang.core.plan.wayangplan.OperatorBase
org.apache.wayang.core.plan.wayangplan.BinaryToUnaryOperator<InputType0,InputType1,Tuple2<Iterable<InputType0>,Iterable<InputType1>>>
org.apache.wayang.basic.operators.CoGroupOperator<InputType0,InputType1,TypeKey>
org.apache.wayang.flink.operators.FlinkCoGroupOperator<InputType0,InputType1,TypeKey>
- All Implemented Interfaces:
Serializable,ActualOperator,ElementaryOperator,ExecutionOperator,Operator,FlinkExecutionOperator
public class FlinkCoGroupOperator<InputType0,InputType1,TypeKey>
extends CoGroupOperator<InputType0,InputType1,TypeKey>
implements FlinkExecutionOperator
Flink implementation of the
CoGroupOperator.- 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.CoGroupOperator
keyDescriptor0, keyDescriptor1Fields 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
ConstructorsConstructorDescriptionFlinkCoGroupOperator(FunctionDescriptor.SerializableFunction<InputType0, TypeKey> keyExtractor0, FunctionDescriptor.SerializableFunction<InputType1, TypeKey> keyExtractor1, Class<InputType0> input0Class, Class<InputType1> input1Class, Class<TypeKey> keyClass) FlinkCoGroupOperator(TransformationDescriptor<InputType0, TypeKey> keyDescriptor0, TransformationDescriptor<InputType1, TypeKey> keyDescriptor1) FlinkCoGroupOperator(TransformationDescriptor<InputType0, TypeKey> keyDescriptor0, TransformationDescriptor<InputType1, TypeKey> keyDescriptor1, DataSetType<InputType0> inputType0, DataSetType<InputType1> inputType1) -
Method Summary
Modifier and TypeMethodDescriptionbooleanTell whether this instances is a Flink action.protected ExecutionOperatorevaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext) getSupportedInputChannels(int index) getSupportedOutputChannels(int index) Display the supportedChannels for a certainOutputSlot.Methods inherited from class org.apache.wayang.basic.operators.CoGroupOperator
createCardinalityEstimator, getKeyDescriptor0, getKeyDescriptor1Methods inherited from class org.apache.wayang.core.plan.wayangplan.BinaryToUnaryOperator
getInputType0, getInputType1, 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, createLoadProfileEstimator, createOutputChannelInstances, getLimitBaseKey, getLoadProfileEstimatorConfigurationKey, getLoadProfileEstimatorConfigurationKeys, getOriginal, getOutputChannelDescriptor, isFilteredMethods inherited from interface org.apache.wayang.flink.operators.FlinkExecutionOperator
getBroadCastFunction, getPlatformMethods 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, setOutput
-
Constructor Details
-
FlinkCoGroupOperator
public FlinkCoGroupOperator(FunctionDescriptor.SerializableFunction<InputType0, TypeKey> keyExtractor0, FunctionDescriptor.SerializableFunction<InputType1, TypeKey> keyExtractor1, Class<InputType0> input0Class, Class<InputType1> input1Class, Class<TypeKey> keyClass) -
FlinkCoGroupOperator
public FlinkCoGroupOperator(TransformationDescriptor<InputType0, TypeKey> keyDescriptor0, TransformationDescriptor<InputType1, TypeKey> keyDescriptor1) -
FlinkCoGroupOperator
public FlinkCoGroupOperator(TransformationDescriptor<InputType0, TypeKey> keyDescriptor0, TransformationDescriptor<InputType1, TypeKey> keyDescriptor1, DataSetType<InputType0> inputType0, DataSetType<InputType1> inputType1) -
FlinkCoGroupOperator
- See Also:
-
-
Method Details
-
evaluate
public Tuple<Collection<ExecutionLineageNode>,Collection<ChannelInstance>> evaluate(ChannelInstance[] inputs, ChannelInstance[] outputs, FlinkExecutor flinkExecutor, OptimizationContext.OperatorContext operatorContext) - Specified by:
evaluatein interfaceFlinkExecutionOperator
-
createCopy
- Overrides:
createCopyin classOperatorBase
-
getLoadProfileEstimatorConfigurationTypeKey
-
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: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
-