@Public public class LocalStreamEnvironment extends StreamExecutionEnvironment
When this environment is instantiated, it uses a default parallelism of 1
. The default
parallelism can be set via StreamExecutionEnvironment.setParallelism(int)
.
cacheFile, DEFAULT_JOB_NAME, isChainingEnabled, transformations
Constructor and Description |
---|
LocalStreamEnvironment()
Creates a new mini cluster stream environment that uses the default configuration.
|
LocalStreamEnvironment(Configuration configuration)
Creates a new mini cluster stream environment that configures its local executor with the given configuration.
|
Modifier and Type | Method and Description |
---|---|
JobExecutionResult |
execute(String jobName)
Executes the JobGraph of the on a mini cluster of CLusterUtil with a user
specified name.
|
protected Configuration |
getConfiguration() |
addDefaultKryoSerializer, addDefaultKryoSerializer, addOperator, addSource, addSource, addSource, addSource, clean, createInput, createInput, createLocalEnvironment, createLocalEnvironment, createLocalEnvironment, createLocalEnvironmentWithWebUI, createRemoteEnvironment, createRemoteEnvironment, createRemoteEnvironment, disableOperatorChaining, enableCheckpointing, enableCheckpointing, enableCheckpointing, enableCheckpointing, execute, fromCollection, fromCollection, fromCollection, fromCollection, fromElements, fromElements, fromParallelCollection, fromParallelCollection, generateSequence, getBufferTimeout, getCachedFiles, getCheckpointConfig, getCheckpointingMode, getCheckpointInterval, getConfig, getDefaultLocalParallelism, getExecutionEnvironment, getExecutionPlan, getMaxParallelism, getNumberOfExecutionRetries, getParallelism, getRestartStrategy, getStateBackend, getStreamGraph, getStreamTimeCharacteristic, initializeContextEnvironment, isChainingEnabled, isForceCheckpointing, readFile, readFile, readFile, readFile, readFileStream, readTextFile, readTextFile, registerCachedFile, registerCachedFile, registerType, registerTypeWithKryoSerializer, registerTypeWithKryoSerializer, resetContextEnvironment, setBufferTimeout, setDefaultLocalParallelism, setMaxParallelism, setNumberOfExecutionRetries, setParallelism, setRestartStrategy, setStateBackend, setStateBackend, setStreamTimeCharacteristic, socketTextStream, socketTextStream, socketTextStream, socketTextStream, socketTextStream
public LocalStreamEnvironment()
public LocalStreamEnvironment(@Nonnull Configuration configuration)
configuration
- The configuration used to configure the local executor.protected Configuration getConfiguration()
public JobExecutionResult execute(String jobName) throws Exception
execute
in class StreamExecutionEnvironment
jobName
- name of the jobException
- which occurs during job execution.Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.