Constructor and Description |
---|
StreamNode(Integer id,
String slotSharingGroup,
String coLocationGroup,
StreamOperator<?> operator,
String operatorName,
Class<? extends TaskInvokable> jobVertexClass) |
StreamNode(Integer id,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<?> operatorFactory,
String operatorName,
Class<? extends TaskInvokable> jobVertexClass) |
@VisibleForTesting public StreamNode(Integer id, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, @Nullable StreamOperator<?> operator, String operatorName, Class<? extends TaskInvokable> jobVertexClass)
public StreamNode(Integer id, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, @Nullable StreamOperatorFactory<?> operatorFactory, String operatorName, Class<? extends TaskInvokable> jobVertexClass)
public void addInEdge(StreamEdge inEdge)
public void addOutEdge(StreamEdge outEdge)
public List<StreamEdge> getOutEdges()
public List<StreamEdge> getInEdges()
public int getId()
public int getParallelism()
public void setParallelism(Integer parallelism)
public ResourceSpec getMinResources()
public ResourceSpec getPreferredResources()
public void setResources(ResourceSpec minResources, ResourceSpec preferredResources)
public void setManagedMemoryUseCaseWeights(Map<ManagedMemoryUseCase,Integer> operatorScopeUseCaseWeights, Set<ManagedMemoryUseCase> slotScopeUseCases)
public Map<ManagedMemoryUseCase,Integer> getManagedMemoryOperatorScopeUseCaseWeights()
public Set<ManagedMemoryUseCase> getManagedMemorySlotScopeUseCases()
public long getBufferTimeout()
public void setBufferTimeout(Long bufferTimeout)
@VisibleForTesting public StreamOperator<?> getOperator()
@Nullable public StreamOperatorFactory<?> getOperatorFactory()
public String getOperatorName()
public String getOperatorDescription()
public void setOperatorDescription(String operatorDescription)
public void setSerializersIn(TypeSerializer<?>... typeSerializersIn)
public TypeSerializer<?>[] getTypeSerializersIn()
public TypeSerializer<?> getTypeSerializerOut()
public void setSerializerOut(TypeSerializer<?> typeSerializerOut)
public Class<? extends TaskInvokable> getJobVertexClass()
public InputFormat<?,?> getInputFormat()
public void setInputFormat(InputFormat<?,?> inputFormat)
public OutputFormat<?> getOutputFormat()
public void setOutputFormat(OutputFormat<?> outputFormat)
public boolean isSameSlotSharingGroup(StreamNode downstreamVertex)
public KeySelector<?,?>[] getStatePartitioners()
public void setStatePartitioners(KeySelector<?,?>... statePartitioners)
public TypeSerializer<?> getStateKeySerializer()
public void setStateKeySerializer(TypeSerializer<?> stateKeySerializer)
public String getTransformationUID()
public String getUserHash()
public void setUserHash(String userHash)
public void addInputRequirement(int inputIndex, StreamConfig.InputRequirement inputRequirement)
public Map<Integer,StreamConfig.InputRequirement> getInputRequirements()
public Optional<OperatorCoordinator.Provider> getCoordinatorProvider(String operatorName, OperatorID operatorID)
@Nullable public IntermediateDataSetID getConsumeClusterDatasetId()
public void setConsumeClusterDatasetId(@Nullable IntermediateDataSetID consumeClusterDatasetId)
public boolean isSupportsConcurrentExecutionAttempts()
public void setSupportsConcurrentExecutionAttempts(boolean supportsConcurrentExecutionAttempts)
public boolean isOutputOnlyAfterEndOfStream()
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.