@PublicEvolving public class StreamPlanEnvironment extends StreamExecutionEnvironment
StreamExecutionEnvironment
that is used in the web frontend when generating
a user-inspectable graph of a streaming job.cacheFile, DEFAULT_JOB_NAME, isChainingEnabled, transformations
Modifier | Constructor and Description |
---|---|
protected |
StreamPlanEnvironment(ExecutionEnvironment env) |
Modifier and Type | Method and Description |
---|---|
JobClient |
executeAsync(StreamGraph streamGraph)
Triggers the program execution asynchronously.
|
addDefaultKryoSerializer, addDefaultKryoSerializer, addOperator, addSource, addSource, addSource, addSource, clean, clearJobListeners, configure, createInput, createInput, createLocalEnvironment, createLocalEnvironment, createLocalEnvironment, createLocalEnvironmentWithWebUI, createRemoteEnvironment, createRemoteEnvironment, createRemoteEnvironment, disableOperatorChaining, enableCheckpointing, enableCheckpointing, enableCheckpointing, enableCheckpointing, execute, execute, execute, executeAsync, executeAsync, fromCollection, fromCollection, fromCollection, fromCollection, fromElements, fromElements, fromParallelCollection, fromParallelCollection, generateSequence, getBufferTimeout, getCachedFiles, getCheckpointConfig, getCheckpointingMode, getCheckpointInterval, getConfig, getConfiguration, getDefaultLocalParallelism, getExecutionEnvironment, getExecutionPlan, getJobListeners, getMaxParallelism, getNumberOfExecutionRetries, getParallelism, getRestartStrategy, getStateBackend, getStreamGraph, getStreamGraph, getStreamGraph, getStreamTimeCharacteristic, initializeContextEnvironment, isChainingEnabled, isForceCheckpointing, readFile, readFile, readFile, readFile, readFileStream, readTextFile, readTextFile, registerCachedFile, registerCachedFile, registerJobListener, registerType, registerTypeWithKryoSerializer, registerTypeWithKryoSerializer, resetContextEnvironment, setBufferTimeout, setDefaultLocalParallelism, setMaxParallelism, setNumberOfExecutionRetries, setParallelism, setRestartStrategy, setStateBackend, setStateBackend, setStreamTimeCharacteristic, socketTextStream, socketTextStream, socketTextStream, socketTextStream, socketTextStream
protected StreamPlanEnvironment(ExecutionEnvironment env)
public JobClient executeAsync(StreamGraph streamGraph) throws Exception
StreamExecutionEnvironment
executeAsync
in class StreamExecutionEnvironment
streamGraph
- the stream graph representing the transformationsJobClient
that can be used to communicate with the submitted job, completed on submission succeeded.Exception
- which occurs during job execution.Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.