Modifier and Type | Class and Description |
---|---|
class |
ScalaShellRemoteStreamEnvironment
A
RemoteStreamEnvironment for the Scala shell. |
Modifier and Type | Method and Description |
---|---|
DataStream<Row> |
StreamSQLTestProgram.GeneratorTableSource.getDataStream(StreamExecutionEnvironment execEnv) |
Modifier and Type | Field and Description |
---|---|
protected StreamExecutionEnvironment |
DataStream.environment |
protected StreamExecutionEnvironment |
ConnectedStreams.environment |
Modifier and Type | Method and Description |
---|---|
StreamExecutionEnvironment |
BroadcastStream.getEnvironment() |
StreamExecutionEnvironment |
AllWindowedStream.getExecutionEnvironment() |
StreamExecutionEnvironment |
BroadcastConnectedStream.getExecutionEnvironment() |
StreamExecutionEnvironment |
DataStream.getExecutionEnvironment()
Returns the
StreamExecutionEnvironment that was used to create this
DataStream . |
StreamExecutionEnvironment |
ConnectedStreams.getExecutionEnvironment() |
StreamExecutionEnvironment |
WindowedStream.getExecutionEnvironment() |
Constructor and Description |
---|
BroadcastConnectedStream(StreamExecutionEnvironment env,
DataStream<IN1> input1,
BroadcastStream<IN2> input2,
List<MapStateDescriptor<?,?>> broadcastStateDescriptors) |
BroadcastStream(StreamExecutionEnvironment env,
DataStream<T> input,
MapStateDescriptor<?,?>... broadcastStateDescriptors) |
ConnectedStreams(StreamExecutionEnvironment env,
DataStream<IN1> input1,
DataStream<IN2> input2) |
DataStream(StreamExecutionEnvironment environment,
StreamTransformation<T> transformation)
Create a new
DataStream in the given execution environment with
partitioning set to forward by default. |
DataStreamSource(StreamExecutionEnvironment environment,
TypeInformation<T> outTypeInfo,
StreamSource<T,?> operator,
boolean isParallel,
String sourceName) |
SingleOutputStreamOperator(StreamExecutionEnvironment environment,
StreamTransformation<T> transformation) |
Modifier and Type | Class and Description |
---|---|
class |
LocalStreamEnvironment
The LocalStreamEnvironment is a StreamExecutionEnvironment that runs the program locally,
multi-threaded, in the JVM where the environment is instantiated.
|
class |
RemoteStreamEnvironment
A
StreamExecutionEnvironment for executing on a cluster. |
class |
StreamContextEnvironment
Special
StreamExecutionEnvironment that will be used in cases where the CLI client or
testing utilities create a StreamExecutionEnvironment that should be used when
getExecutionEnvironment() is called. |
class |
StreamPlanEnvironment
A special
StreamExecutionEnvironment that is used in the web frontend when generating
a user-inspectable graph of a streaming job. |
Modifier and Type | Method and Description |
---|---|
StreamExecutionEnvironment |
StreamExecutionEnvironmentFactory.createExecutionEnvironment()
Creates a StreamExecutionEnvironment from this factory.
|
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(Configuration conf)
Creates a
LocalStreamEnvironment for local program execution that also starts the
web monitoring UI. |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfig,
String... jarFiles)
Creates a
RemoteStreamEnvironment . |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createRemoteEnvironment(String host,
int port,
int parallelism,
String... jarFiles)
Creates a
RemoteStreamEnvironment . |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createRemoteEnvironment(String host,
int port,
String... jarFiles)
Creates a
RemoteStreamEnvironment . |
StreamExecutionEnvironment |
StreamExecutionEnvironment.disableOperatorChaining()
Disables operator chaining for streaming operators.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.enableCheckpointing()
Deprecated.
Use
enableCheckpointing(long) instead. |
StreamExecutionEnvironment |
StreamExecutionEnvironment.enableCheckpointing(long interval)
Enables checkpointing for the streaming job.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.enableCheckpointing(long interval,
CheckpointingMode mode)
Enables checkpointing for the streaming job.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.enableCheckpointing(long interval,
CheckpointingMode mode,
boolean force)
Deprecated.
Use
enableCheckpointing(long, CheckpointingMode) instead.
Forcing checkpoints will be removed in the future. |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.getExecutionEnvironment()
Creates an execution environment that represents the context in which the
program is currently executed.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.setBufferTimeout(long timeoutMillis)
Sets the maximum time frequency (milliseconds) for the flushing of the
output buffers.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.setMaxParallelism(int maxParallelism)
Sets the maximum degree of parallelism defined for the program.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.setParallelism(int parallelism)
Sets the parallelism for operations executed through this environment.
|
StreamExecutionEnvironment |
StreamExecutionEnvironment.setStateBackend(AbstractStateBackend backend)
Deprecated.
Use
setStateBackend(StateBackend) instead. |
StreamExecutionEnvironment |
StreamExecutionEnvironment.setStateBackend(StateBackend backend)
Sets the state backend that describes how to store and checkpoint operator state.
|
Modifier and Type | Method and Description |
---|---|
static JobExecutionResult |
RemoteStreamEnvironment.executeRemotely(StreamExecutionEnvironment streamExecutionEnvironment,
List<URL> jarFiles,
String host,
int port,
Configuration clientConfiguration,
List<URL> globalClasspaths,
String jobName,
SavepointRestoreSettings savepointRestoreSettings)
Executes the job remotely.
|
Modifier and Type | Method and Description |
---|---|
StreamExecutionEnvironment |
StreamGraph.getEnvironment() |
Modifier and Type | Method and Description |
---|---|
static StreamGraph |
StreamGraphGenerator.generate(StreamExecutionEnvironment env,
List<StreamTransformation<?>> transformations)
Generates a
StreamGraph by traversing the graph of StreamTransformations
starting from the given transformations. |
Constructor and Description |
---|
StreamGraph(StreamExecutionEnvironment environment) |
StreamNode(StreamExecutionEnvironment env,
Integer id,
String slotSharingGroup,
String coLocationGroup,
StreamOperator<?> operator,
String operatorName,
List<OutputSelector<?>> outputSelector,
Class<? extends AbstractInvokable> jobVertexClass) |
Modifier and Type | Method and Description |
---|---|
DataStream<Row> |
KafkaTableSourceBase.getDataStream(StreamExecutionEnvironment env)
NOTE: This method is for internal use only for defining a TableSource.
|
Modifier and Type | Method and Description |
---|---|
static DataStream<Tuple2<String,Integer>> |
WindowJoinSampleData.GradeSource.getSource(StreamExecutionEnvironment env,
long rate) |
static DataStream<Tuple2<String,Integer>> |
WindowJoinSampleData.SalarySource.getSource(StreamExecutionEnvironment env,
long rate) |
Modifier and Type | Method and Description |
---|---|
static StreamExecutionEnvironment |
KafkaExampleUtil.prepareExecutionEnv(ParameterTool parameterTool) |
Modifier and Type | Method and Description |
---|---|
static void |
DataStreamAllroundTestJobFactory.setupEnvironment(StreamExecutionEnvironment env,
ParameterTool pt) |
Modifier and Type | Class and Description |
---|---|
class |
TestStreamEnvironment
A
StreamExecutionEnvironment that executes its jobs on MiniCluster . |
Modifier and Type | Method and Description |
---|---|
StreamExecutionEnvironment |
ExecutionContext.EnvironmentInstance.getStreamExecutionEnvironment() |
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.