Modifier and Type | Field and Description |
---|---|
static String |
ITERATION_SINK_NAME_PREFIX |
static String |
ITERATION_SOURCE_NAME_PREFIX |
protected Map<Integer,String> |
vertexIDtoBrokerID |
protected Map<Integer,Long> |
vertexIDtoLoopTimeout |
Constructor and Description |
---|
StreamGraph(ExecutionConfig executionConfig,
CheckpointConfig checkpointConfig,
SavepointRestoreSettings savepointRestoreSettings) |
Modifier and Type | Method and Description |
---|---|
<IN1,IN2,OUT> |
addCoOperator(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<OUT> taskOperatorFactory,
TypeInformation<IN1> in1TypeInfo,
TypeInformation<IN2> in2TypeInfo,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
void |
addEdge(Integer upStreamVertexID,
Integer downStreamVertexID,
int typeNumber) |
<IN,OUT> void |
addLegacySource(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<OUT> operatorFactory,
TypeInformation<IN> inTypeInfo,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
<OUT> void |
addMultipleInputOperator(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<OUT> operatorFactory,
List<TypeInformation<?>> inTypeInfos,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
protected StreamNode |
addNode(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
Class<? extends AbstractInvokable> vertexClass,
StreamOperatorFactory<?> operatorFactory,
String operatorName) |
<IN,OUT> void |
addOperator(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<OUT> operatorFactory,
TypeInformation<IN> inTypeInfo,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
<IN,OUT> void |
addSink(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
StreamOperatorFactory<OUT> operatorFactory,
TypeInformation<IN> inTypeInfo,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
<IN,OUT> void |
addSource(Integer vertexID,
String slotSharingGroup,
String coLocationGroup,
SourceOperatorFactory<OUT> operatorFactory,
TypeInformation<IN> inTypeInfo,
TypeInformation<OUT> outTypeInfo,
String operatorName) |
void |
addVirtualPartitionNode(Integer originalId,
Integer virtualId,
StreamPartitioner<?> partitioner,
ShuffleMode shuffleMode)
Adds a new virtual node that is used to connect a downstream vertex to an input with a
certain partitioning.
|
void |
addVirtualSideOutputNode(Integer originalId,
Integer virtualId,
OutputTag outputTag)
Adds a new virtual node that is used to connect a downstream vertex to only the outputs with
the selected side-output
OutputTag . |
void |
clear()
Remove all registered nodes etc.
|
Tuple2<StreamNode,StreamNode> |
createIterationSourceAndSink(int loopId,
int sourceId,
int sinkId,
long timeout,
int parallelism,
int maxParallelism,
ResourceSpec minResources,
ResourceSpec preferredResources) |
Set<Tuple2<Integer,StreamOperatorFactory<?>>> |
getAllOperatorFactory() |
String |
getBrokerID(Integer vertexID) |
CheckpointConfig |
getCheckpointConfig() |
CheckpointStorage |
getCheckpointStorage() |
ExecutionConfig |
getExecutionConfig() |
GlobalDataExchangeMode |
getGlobalDataExchangeMode() |
Set<Tuple2<StreamNode,StreamNode>> |
getIterationSourceSinkPairs() |
JobGraph |
getJobGraph()
|
JobGraph |
getJobGraph(JobID jobID)
|
String |
getJobName() |
JobType |
getJobType() |
long |
getLoopTimeout(Integer vertexID) |
Path |
getSavepointDirectory() |
SavepointRestoreSettings |
getSavepointRestoreSettings() |
Collection<Integer> |
getSinkIDs() |
String |
getSlotSharingGroup(Integer id)
Determines the slot sharing group of an operation across virtual nodes.
|
Optional<ResourceProfile> |
getSlotSharingGroupResource(String groupId) |
Collection<Integer> |
getSourceIDs() |
StreamNode |
getSourceVertex(StreamEdge edge) |
StateBackend |
getStateBackend() |
List<StreamEdge> |
getStreamEdges(int sourceId) |
List<StreamEdge> |
getStreamEdges(int sourceId,
int targetId) |
List<StreamEdge> |
getStreamEdgesOrThrow(int sourceId,
int targetId)
Deprecated.
|
String |
getStreamingPlanAsJSON() |
StreamNode |
getStreamNode(Integer vertexID) |
Collection<StreamNode> |
getStreamNodes() |
StreamNode |
getTargetVertex(StreamEdge edge) |
TimeCharacteristic |
getTimeCharacteristic() |
InternalTimeServiceManager.Provider |
getTimerServiceProvider() |
Collection<Tuple2<String,DistributedCache.DistributedCacheEntry>> |
getUserArtifacts() |
protected Collection<? extends Integer> |
getVertexIDs() |
boolean |
isAllVerticesInSameSlotSharingGroupByDefault()
Gets whether to put all vertices into the same slot sharing group by default.
|
boolean |
isChainingEnabled() |
boolean |
isIterative() |
void |
setAllVerticesInSameSlotSharingGroupByDefault(boolean allVerticesInSameSlotSharingGroupByDefault)
Set whether to put all vertices into the same slot sharing group by default.
|
void |
setBufferTimeout(Integer vertexID,
long bufferTimeout) |
void |
setChaining(boolean chaining) |
void |
setCheckpointStorage(CheckpointStorage checkpointStorage) |
void |
setGlobalDataExchangeMode(GlobalDataExchangeMode globalDataExchangeMode) |
void |
setInputFormat(Integer vertexID,
InputFormat<?,?> inputFormat) |
void |
setJobName(String jobName) |
void |
setJobType(JobType jobType) |
void |
setManagedMemoryUseCaseWeights(int vertexID,
Map<ManagedMemoryUseCase,Integer> operatorScopeUseCaseWeights,
Set<ManagedMemoryUseCase> slotScopeUseCases) |
void |
setMaxParallelism(int vertexID,
int maxParallelism) |
void |
setMultipleInputStateKey(Integer vertexID,
List<KeySelector<?,?>> keySelectors,
TypeSerializer<?> keySerializer) |
void |
setOneInputStateKey(Integer vertexID,
KeySelector<?,?> keySelector,
TypeSerializer<?> keySerializer) |
void |
setOutputFormat(Integer vertexID,
OutputFormat<?> outputFormat) |
<OUT> void |
setOutType(Integer vertexID,
TypeInformation<OUT> outType) |
void |
setParallelism(Integer vertexID,
int parallelism) |
void |
setResources(int vertexID,
ResourceSpec minResources,
ResourceSpec preferredResources) |
void |
setSavepointDirectory(Path savepointDir) |
void |
setSavepointRestoreSettings(SavepointRestoreSettings savepointRestoreSettings) |
void |
setSerializers(Integer vertexID,
TypeSerializer<?> in1,
TypeSerializer<?> in2,
TypeSerializer<?> out) |
void |
setSerializersFrom(Integer from,
Integer to) |
void |
setSlotSharingGroupResource(Map<String,ResourceProfile> slotSharingGroupResources) |
void |
setStateBackend(StateBackend backend) |
void |
setTimeCharacteristic(TimeCharacteristic timeCharacteristic) |
void |
setTimerServiceProvider(InternalTimeServiceManager.Provider timerServiceProvider) |
void |
setTransformationUID(Integer nodeId,
String transformationId) |
void |
setTwoInputStateKey(Integer vertexID,
KeySelector<?,?> keySelector1,
KeySelector<?,?> keySelector2,
TypeSerializer<?> keySerializer) |
void |
setUserArtifacts(Collection<Tuple2<String,DistributedCache.DistributedCacheEntry>> userArtifacts) |
public static final String ITERATION_SOURCE_NAME_PREFIX
public static final String ITERATION_SINK_NAME_PREFIX
public StreamGraph(ExecutionConfig executionConfig, CheckpointConfig checkpointConfig, SavepointRestoreSettings savepointRestoreSettings)
public void clear()
public ExecutionConfig getExecutionConfig()
public CheckpointConfig getCheckpointConfig()
public void setSavepointRestoreSettings(SavepointRestoreSettings savepointRestoreSettings)
public SavepointRestoreSettings getSavepointRestoreSettings()
public String getJobName()
public void setJobName(String jobName)
public void setChaining(boolean chaining)
public void setStateBackend(StateBackend backend)
public StateBackend getStateBackend()
public void setCheckpointStorage(CheckpointStorage checkpointStorage)
public CheckpointStorage getCheckpointStorage()
public void setSavepointDirectory(Path savepointDir)
public Path getSavepointDirectory()
public InternalTimeServiceManager.Provider getTimerServiceProvider()
public void setTimerServiceProvider(InternalTimeServiceManager.Provider timerServiceProvider)
public Collection<Tuple2<String,DistributedCache.DistributedCacheEntry>> getUserArtifacts()
public void setUserArtifacts(Collection<Tuple2<String,DistributedCache.DistributedCacheEntry>> userArtifacts)
public TimeCharacteristic getTimeCharacteristic()
public void setTimeCharacteristic(TimeCharacteristic timeCharacteristic)
public GlobalDataExchangeMode getGlobalDataExchangeMode()
public void setGlobalDataExchangeMode(GlobalDataExchangeMode globalDataExchangeMode)
public void setSlotSharingGroupResource(Map<String,ResourceProfile> slotSharingGroupResources)
public Optional<ResourceProfile> getSlotSharingGroupResource(String groupId)
public void setAllVerticesInSameSlotSharingGroupByDefault(boolean allVerticesInSameSlotSharingGroupByDefault)
allVerticesInSameSlotSharingGroupByDefault
- indicates whether to put all vertices into
the same slot sharing group by default.public boolean isAllVerticesInSameSlotSharingGroupByDefault()
public boolean isChainingEnabled()
public boolean isIterative()
public <IN,OUT> void addSource(Integer vertexID, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, SourceOperatorFactory<OUT> operatorFactory, TypeInformation<IN> inTypeInfo, TypeInformation<OUT> outTypeInfo, String operatorName)
public <IN,OUT> void addLegacySource(Integer vertexID, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, StreamOperatorFactory<OUT> operatorFactory, TypeInformation<IN> inTypeInfo, TypeInformation<OUT> outTypeInfo, String operatorName)
public <IN,OUT> void addSink(Integer vertexID, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, StreamOperatorFactory<OUT> operatorFactory, TypeInformation<IN> inTypeInfo, TypeInformation<OUT> outTypeInfo, String operatorName)
public <IN,OUT> void addOperator(Integer vertexID, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, StreamOperatorFactory<OUT> operatorFactory, TypeInformation<IN> inTypeInfo, TypeInformation<OUT> outTypeInfo, String operatorName)
public <IN1,IN2,OUT> void addCoOperator(Integer vertexID, String slotSharingGroup, @Nullable String coLocationGroup, StreamOperatorFactory<OUT> taskOperatorFactory, TypeInformation<IN1> in1TypeInfo, TypeInformation<IN2> in2TypeInfo, TypeInformation<OUT> outTypeInfo, String operatorName)
public <OUT> void addMultipleInputOperator(Integer vertexID, String slotSharingGroup, @Nullable String coLocationGroup, StreamOperatorFactory<OUT> operatorFactory, List<TypeInformation<?>> inTypeInfos, TypeInformation<OUT> outTypeInfo, String operatorName)
protected StreamNode addNode(Integer vertexID, @Nullable String slotSharingGroup, @Nullable String coLocationGroup, Class<? extends AbstractInvokable> vertexClass, StreamOperatorFactory<?> operatorFactory, String operatorName)
public void addVirtualSideOutputNode(Integer originalId, Integer virtualId, OutputTag outputTag)
OutputTag
.originalId
- ID of the node that should be connected to.virtualId
- ID of the virtual node.outputTag
- The selected side-output OutputTag
.public void addVirtualPartitionNode(Integer originalId, Integer virtualId, StreamPartitioner<?> partitioner, ShuffleMode shuffleMode)
When adding an edge from the virtual node to a downstream node the connection will be made to the original node, but with the partitioning given here.
originalId
- ID of the node that should be connected to.virtualId
- ID of the virtual node.partitioner
- The partitionerpublic String getSlotSharingGroup(Integer id)
public void setParallelism(Integer vertexID, int parallelism)
public void setMaxParallelism(int vertexID, int maxParallelism)
public void setResources(int vertexID, ResourceSpec minResources, ResourceSpec preferredResources)
public void setManagedMemoryUseCaseWeights(int vertexID, Map<ManagedMemoryUseCase,Integer> operatorScopeUseCaseWeights, Set<ManagedMemoryUseCase> slotScopeUseCases)
public void setOneInputStateKey(Integer vertexID, KeySelector<?,?> keySelector, TypeSerializer<?> keySerializer)
public void setTwoInputStateKey(Integer vertexID, KeySelector<?,?> keySelector1, KeySelector<?,?> keySelector2, TypeSerializer<?> keySerializer)
public void setMultipleInputStateKey(Integer vertexID, List<KeySelector<?,?>> keySelectors, TypeSerializer<?> keySerializer)
public void setBufferTimeout(Integer vertexID, long bufferTimeout)
public void setSerializers(Integer vertexID, TypeSerializer<?> in1, TypeSerializer<?> in2, TypeSerializer<?> out)
public <OUT> void setOutType(Integer vertexID, TypeInformation<OUT> outType)
public void setInputFormat(Integer vertexID, InputFormat<?,?> inputFormat)
public void setOutputFormat(Integer vertexID, OutputFormat<?> outputFormat)
public StreamNode getStreamNode(Integer vertexID)
protected Collection<? extends Integer> getVertexIDs()
@VisibleForTesting public List<StreamEdge> getStreamEdges(int sourceId)
@VisibleForTesting public List<StreamEdge> getStreamEdges(int sourceId, int targetId)
@VisibleForTesting @Deprecated public List<StreamEdge> getStreamEdgesOrThrow(int sourceId, int targetId)
public Collection<Integer> getSourceIDs()
public Collection<Integer> getSinkIDs()
public Collection<StreamNode> getStreamNodes()
public Set<Tuple2<Integer,StreamOperatorFactory<?>>> getAllOperatorFactory()
public long getLoopTimeout(Integer vertexID)
public Tuple2<StreamNode,StreamNode> createIterationSourceAndSink(int loopId, int sourceId, int sinkId, long timeout, int parallelism, int maxParallelism, ResourceSpec minResources, ResourceSpec preferredResources)
public Set<Tuple2<StreamNode,StreamNode>> getIterationSourceSinkPairs()
public StreamNode getSourceVertex(StreamEdge edge)
public StreamNode getTargetVertex(StreamEdge edge)
public JobGraph getJobGraph()
public String getStreamingPlanAsJSON()
public void setJobType(JobType jobType)
public JobType getJobType()
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.