Modifier and Type | Method and Description |
---|---|
void |
TableInputFormat.configure(Configuration parameters)
Creates a
Scan object and opens the HTable connection. |
Modifier and Type | Method and Description |
---|---|
static PlanExecutor |
PlanExecutor.createLocalExecutor(Configuration configuration)
Creates an executor that runs the plan locally in a multi-threaded environment.
|
static PlanExecutor |
PlanExecutor.createRemoteExecutor(String hostname,
int port,
Configuration clientConfiguration,
List<URL> jarFiles,
List<URL> globalClasspaths)
Creates an executor that runs the plan on a remote environment.
|
Modifier and Type | Method and Description |
---|---|
static Set<Map.Entry<String,DistributedCache.DistributedCacheEntry>> |
DistributedCache.readFileInfoFromConfig(Configuration conf) |
static void |
DistributedCache.writeFileInfoToConfig(String name,
DistributedCache.DistributedCacheEntry e,
Configuration conf) |
Modifier and Type | Method and Description |
---|---|
void |
RichFunction.open(Configuration parameters)
Initialization method for the function.
|
void |
AbstractRichFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
static void |
FunctionUtils.openFunction(Function function,
Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
ReplicatingInputFormat.configure(Configuration parameters) |
void |
OutputFormat.configure(Configuration parameters)
Configures this output format.
|
void |
InputFormat.configure(Configuration parameters)
Configures this input format.
|
void |
GenericInputFormat.configure(Configuration parameters) |
void |
FileOutputFormat.configure(Configuration parameters) |
void |
FileInputFormat.configure(Configuration parameters)
Configures the file input format by reading the file path from the configuration.
|
void |
DelimitedInputFormat.configure(Configuration parameters)
Configures this input format by reading the path to the file from the configuration andge the string that
defines the record delimiter.
|
void |
BinaryOutputFormat.configure(Configuration parameters) |
void |
BinaryInputFormat.configure(Configuration parameters) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
Operator.parameters |
Modifier and Type | Method and Description |
---|---|
Configuration |
Operator.getParameters()
Gets the stub parameters of this contract.
|
Modifier and Type | Method and Description |
---|---|
void |
BulkIterationBase.TerminationCriterionMapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
TypeSerializerFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
TypeComparatorFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
TypeSerializerFactory.writeParametersToConfig(Configuration config) |
void |
TypeComparatorFactory.writeParametersToConfig(Configuration config) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
RemoteEnvironment.clientConfiguration
The configuration used by the client that connects to the cluster
|
Modifier and Type | Method and Description |
---|---|
void |
Utils.CountHelper.configure(Configuration parameters) |
void |
Utils.CollectHelper.configure(Configuration parameters) |
void |
Utils.ChecksumHashCodeHelper.configure(Configuration parameters) |
static LocalEnvironment |
ExecutionEnvironment.createLocalEnvironment(Configuration customConfiguration)
Creates a
LocalEnvironment . |
static ExecutionEnvironment |
ExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
String... jarFiles)
Creates a
RemoteEnvironment . |
Constructor and Description |
---|
LocalEnvironment(Configuration config)
Creates a new local environment that configures its local executor with the given configuration.
|
RemoteEnvironment(String host,
int port,
Configuration clientConfig,
String[] jarFiles)
Creates a new RemoteEnvironment that points to the master (JobManager) described by the
given host name and port.
|
RemoteEnvironment(String host,
int port,
Configuration clientConfig,
String[] jarFiles,
URL[] globalClasspaths)
Creates a new RemoteEnvironment that points to the master (JobManager) described by the
given host name and port.
|
ScalaShellRemoteEnvironment(String host,
int port,
FlinkILoop flinkILoop,
Configuration clientConfig,
String... jarFiles)
Creates new ScalaShellRemoteEnvironment that has a reference to the FlinkILoop
|
ScalaShellRemoteStreamEnvironment(String host,
int port,
FlinkILoop flinkILoop,
Configuration configuration,
String... jarFiles)
Creates a new RemoteStreamEnvironment that points to the master
(JobManager) described by the given host name and port.
|
Modifier and Type | Method and Description |
---|---|
void |
HadoopOutputFormatBase.configure(Configuration parameters) |
void |
HadoopInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopOutputFormatBase.configure(Configuration parameters) |
void |
HadoopInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
TextValueInputFormat.configure(Configuration parameters) |
void |
TextInputFormat.configure(Configuration parameters) |
void |
PrintingOutputFormat.configure(Configuration parameters) |
void |
LocalCollectionOutputFormat.configure(Configuration parameters) |
void |
DiscardingOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
JDBCInputFormat.configure(Configuration parameters) |
void |
JDBCOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
Configuration |
UdfOperator.getParameters()
Gets the configuration parameters that will be passed to the UDF's open method
AbstractRichFunction.open(Configuration) . |
Configuration |
TwoInputUdfOperator.getParameters() |
Configuration |
SingleInputUdfOperator.getParameters() |
Configuration |
DataSource.getParameters() |
Configuration |
DataSink.getParameters() |
Modifier and Type | Method and Description |
---|---|
void |
AggregateOperator.AggregatingUdf.open(Configuration parameters) |
O |
UdfOperator.withParameters(Configuration parameters)
Sets the configuration parameters for the UDF.
|
O |
TwoInputUdfOperator.withParameters(Configuration parameters) |
O |
SingleInputUdfOperator.withParameters(Configuration parameters) |
DataSource<OUT> |
DataSource.withParameters(Configuration parameters)
Pass a configuration to the InputFormat
|
DataSink<T> |
DataSink.withParameters(Configuration parameters)
Pass a configuration to the OutputFormat
|
Modifier and Type | Method and Description |
---|---|
void |
WrappingFunction.open(Configuration parameters) |
void |
RichCombineToGroupCombineWrapper.open(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
RuntimeSerializerFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeComparatorFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeSerializerFactory.writeParametersToConfig(Configuration config) |
void |
RuntimeComparatorFactory.writeParametersToConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ParameterTool.getConfiguration()
Returns a
Configuration object from this ParameterTool |
Modifier and Type | Method and Description |
---|---|
Configuration |
FlinkILoop.clientConfig() |
Modifier and Type | Method and Description |
---|---|
static ExecutionEnvironment |
ExecutionEnvironment.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
ExecutionEnvironment |
ExecutionEnvironment$.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
static ExecutionEnvironment |
ExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
scala.collection.Seq<String> jarFiles)
Creates a remote execution environment.
|
ExecutionEnvironment |
ExecutionEnvironment$.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
scala.collection.Seq<String> jarFiles)
Creates a remote execution environment.
|
DataSet<T> |
DataSet.withParameters(Configuration parameters) |
Constructor and Description |
---|
FlinkILoop(String host,
int port,
Configuration clientConfig,
BufferedReader in0,
PrintWriter out) |
FlinkILoop(String host,
int port,
Configuration clientConfig,
scala.Option<String[]> externalJars) |
FlinkILoop(String host,
int port,
Configuration clientConfig,
scala.Option<String[]> externalJars,
BufferedReader in0,
PrintWriter out) |
FlinkILoop(String host,
int port,
Configuration clientConfig,
scala.Option<String[]> externalJars,
scala.Option<BufferedReader> in0,
PrintWriter out0) |
Modifier and Type | Method and Description |
---|---|
void |
ScalaAggregateOperator.AggregatingUdf.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
MapRunner.open(Configuration parameters) |
void |
FlatMapRunner.open(Configuration parameters) |
void |
FlatJoinRunner.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
AggregateMapFunction.open(Configuration config) |
void |
AggregateReduceGroupFunction.open(Configuration config) |
void |
AggregateReduceCombineFunction.open(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
CassandraOutputFormat.configure(Configuration parameters) |
void |
CassandraInputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
Configuration |
CliFrontend.getConfiguration()
Getter which returns a copy of the associated configuration
|
Modifier and Type | Method and Description |
---|---|
static void |
CliFrontend.setJobManagerAddressInConfig(Configuration config,
InetSocketAddress address)
Writes the given job manager address to the associated configuration object
|
Constructor and Description |
---|
LocalExecutor(Configuration conf) |
RemoteExecutor(InetSocketAddress inet,
Configuration clientConfiguration,
List<URL> jarFiles,
List<URL> globalClasspaths) |
RemoteExecutor(String hostport,
Configuration clientConfiguration,
URL jarFile) |
RemoteExecutor(String hostname,
int port,
Configuration clientConfiguration) |
RemoteExecutor(String hostname,
int port,
Configuration clientConfiguration,
List<URL> jarFiles,
List<URL> globalClasspaths) |
RemoteExecutor(String hostname,
int port,
Configuration clientConfiguration,
URL jarFile) |
Modifier and Type | Method and Description |
---|---|
StandaloneClusterClient |
DefaultCLI.createCluster(String applicationName,
org.apache.commons.cli.CommandLine commandLine,
Configuration config) |
ClusterType |
CustomCommandLine.createCluster(String applicationName,
org.apache.commons.cli.CommandLine commandLine,
Configuration config)
Creates the client for the cluster
|
boolean |
DefaultCLI.isActive(org.apache.commons.cli.CommandLine commandLine,
Configuration configuration) |
boolean |
CustomCommandLine.isActive(org.apache.commons.cli.CommandLine commandLine,
Configuration configuration)
Signals whether the custom command-line wants to execute or not
|
StandaloneClusterClient |
DefaultCLI.retrieveCluster(org.apache.commons.cli.CommandLine commandLine,
Configuration config) |
ClusterType |
CustomCommandLine.retrieveCluster(org.apache.commons.cli.CommandLine commandLine,
Configuration config)
Retrieves a client for a running cluster
|
Constructor and Description |
---|
StandaloneClusterDescriptor(Configuration config) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
ClusterClient.flinkConfig
Configuration of the client
|
Modifier and Type | Method and Description |
---|---|
Configuration |
ClusterClient.getFlinkConfiguration()
Return the Flink configuration object
|
Constructor and Description |
---|
ClusterClient(Configuration flinkConfig)
Creates a instance that submits the programs to the JobManager defined in the
configuration.
|
StandaloneClusterClient(Configuration config) |
Modifier and Type | Class and Description |
---|---|
class |
DelegatingConfiguration
A configuration that manages a subset of keys with a common prefix from a given configuration.
|
class |
UnmodifiableConfiguration
Unmodifiable version of the Configuration class.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
Configuration.clone() |
static Configuration |
GlobalConfiguration.getConfiguration()
Gets a
Configuration object with the values of this GlobalConfiguration |
Modifier and Type | Method and Description |
---|---|
void |
UnmodifiableConfiguration.addAll(Configuration other) |
void |
DelegatingConfiguration.addAll(Configuration other) |
void |
Configuration.addAll(Configuration other) |
void |
UnmodifiableConfiguration.addAll(Configuration other,
String prefix) |
void |
DelegatingConfiguration.addAll(Configuration other,
String prefix) |
void |
Configuration.addAll(Configuration other,
String prefix)
Adds all entries from the given configuration into this configuration.
|
static void |
GlobalConfiguration.includeConfiguration(Configuration conf)
Merges the given
Configuration object into the global
configuration. |
Constructor and Description |
---|
Configuration(Configuration other)
Creates a new configuration with the copy of the given configuration.
|
DelegatingConfiguration(Configuration backingConfig,
String prefix)
Creates a new delegating configuration which stores its key/value pairs in the given
configuration using the specifies key prefix.
|
UnmodifiableConfiguration(Configuration config)
Creates a new UnmodifiableConfiguration, which holds a copy of the given configuration
that cannot be altered.
|
Modifier and Type | Method and Description |
---|---|
static void |
FileSystem.setDefaultScheme(Configuration config)
Sets the default filesystem scheme based on the user-specified configuration parameter
fs.default-scheme . |
Modifier and Type | Method and Description |
---|---|
void |
KMeans.SelectNearestCenter.open(Configuration parameters)
Reads the centroid values from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
FileCopyTaskInputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
LinearRegression.SubUpdate.open(Configuration parameters)
Reads the parameters from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
EmptyFieldsCountAccumulator.EmptyFieldFilter.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
KMeans.SelectNearestCenter.open(Configuration parameters)
Reads the centroid values from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
LinearRegression.SubUpdate.open(Configuration parameters)
Reads the parameters from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
ExampleUtils.PrintingOutputFormatWithMessage.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopReduceFunction.open(Configuration parameters) |
void |
HadoopMapFunction.open(Configuration parameters) |
void |
HadoopReduceCombineFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HCatInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HCatInputFormat.configure(Configuration parameters) |
Constructor and Description |
---|
Optimizer(Configuration config)
Creates a new optimizer instance.
|
Optimizer(CostEstimator estimator,
Configuration config)
Creates a new optimizer instance.
|
Optimizer(DataStatistics stats,
Configuration config)
Creates a new optimizer instance that uses the statistics object to determine properties about the input.
|
Optimizer(DataStatistics stats,
CostEstimator estimator,
Configuration config)
Creates a new optimizer instance that uses the statistics object to determine properties about the input.
|
Constructor and Description |
---|
JobGraphGenerator(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
PythonMapPartition.open(Configuration config)
Opens this function.
|
void |
PythonCoGroup.open(Configuration config)
Opens this function.
|
Modifier and Type | Method and Description |
---|---|
void |
PythonStreamer.sendBroadCastVariables(Configuration config)
Sends all broadcast-variables encoded in the configuration to the external process.
|
Modifier and Type | Method and Description |
---|---|
akka.actor.ActorSystem |
AkkaUtils$.createActorSystem(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an actor system.
|
static akka.actor.ActorSystem |
AkkaUtils.createActorSystem(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an actor system.
|
akka.actor.ActorSystem |
AkkaUtils$.createLocalActorSystem(Configuration configuration)
Creates a local actor system without remoting.
|
static akka.actor.ActorSystem |
AkkaUtils.createLocalActorSystem(Configuration configuration)
Creates a local actor system without remoting.
|
com.typesafe.config.Config |
AkkaUtils$.getAkkaConfig(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an akka config with the provided configuration values.
|
static com.typesafe.config.Config |
AkkaUtils.getAkkaConfig(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an akka config with the provided configuration values.
|
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getClientTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getClientTimeout(Configuration config) |
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getLookupTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getLookupTimeout(Configuration config) |
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getTimeout(Configuration config) |
Constructor and Description |
---|
BlobCache(InetSocketAddress serverAddress,
Configuration configuration) |
BlobServer(Configuration config)
Instantiates a new BLOB server and binds it to a free network port.
|
Constructor and Description |
---|
ZooKeeperCheckpointRecoveryFactory(org.apache.curator.framework.CuratorFramework client,
Configuration config,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
static SavepointStore |
SavepointStoreFactory.createFromConfig(Configuration config)
Creates a
SavepointStore from the specified Configuration. |
Modifier and Type | Method and Description |
---|---|
static akka.actor.ActorSystem |
JobClient.startJobClientActorSystem(Configuration config) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
FlinkResourceManager.config
The Flink configuration object
|
Modifier and Type | Method and Description |
---|---|
static Configuration |
BootstrapTools.generateTaskManagerConfiguration(Configuration baseConfig,
String jobManagerHostname,
int jobManagerPort,
int numSlots,
scala.concurrent.duration.FiniteDuration registrationTimeout)
Generate a task manager configuration.
|
Modifier and Type | Method and Description |
---|---|
static ContaineredTaskManagerParameters |
ContaineredTaskManagerParameters.create(Configuration config,
long containerMemoryMB,
int numSlots)
Computes the parameters to be used to start a TaskManager Java process.
|
static Configuration |
BootstrapTools.generateTaskManagerConfiguration(Configuration baseConfig,
String jobManagerHostname,
int jobManagerPort,
int numSlots,
scala.concurrent.duration.FiniteDuration registrationTimeout)
Generate a task manager configuration.
|
static String |
BootstrapTools.getTaskManagerShellCommand(Configuration flinkConfig,
ContaineredTaskManagerParameters tmParams,
String configDirectory,
String logDirectory,
boolean hasLogback,
boolean hasLog4j,
Class<?> mainClass)
Generates the shell command to start a task manager.
|
static akka.actor.ActorSystem |
BootstrapTools.startActorSystem(Configuration configuration,
String listeningAddress,
int listeningPort,
org.slf4j.Logger logger)
Starts an Actor System at a specific port.
|
static akka.actor.ActorSystem |
BootstrapTools.startActorSystem(Configuration configuration,
String listeningAddress,
String portRangeDefinition,
org.slf4j.Logger logger)
Starts an ActorSystem with the given configuration listening at the address/ports.
|
static akka.actor.ActorRef |
FlinkResourceManager.startResourceManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
LeaderRetrievalService leaderRetriever,
Class<? extends FlinkResourceManager<?>> resourceManagerClass)
Starts the resource manager actors.
|
static akka.actor.ActorRef |
FlinkResourceManager.startResourceManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
LeaderRetrievalService leaderRetriever,
Class<? extends FlinkResourceManager<?>> resourceManagerClass,
String resourceManagerActorName)
Starts the resource manager actors.
|
static WebMonitor |
BootstrapTools.startWebMonitorIfConfigured(Configuration config,
akka.actor.ActorSystem actorSystem,
akka.actor.ActorRef jobManager,
org.slf4j.Logger logger)
Starts the web frontend.
|
static void |
BootstrapTools.substituteDeprecatedConfigKey(Configuration config,
String deprecated,
String designated)
Sets the value of a new config key to the value of a deprecated config key.
|
static void |
BootstrapTools.substituteDeprecatedConfigPrefix(Configuration config,
String deprecatedPrefix,
String designatedPrefix)
Sets the value of of a new config key to the value of a deprecated config key.
|
static void |
BootstrapTools.writeConfiguration(Configuration cfg,
File file)
Writes a Flink YAML config file from a Flink Configuration object.
|
Constructor and Description |
---|
FlinkResourceManager(int numInitialTaskManagers,
Configuration flinkConfig,
LeaderRetrievalService leaderRetriever)
Creates a AbstractFrameworkMaster actor.
|
Constructor and Description |
---|
StandaloneResourceManager(Configuration flinkConfig,
LeaderRetrievalService leaderRetriever) |
Modifier and Type | Method and Description |
---|---|
Configuration |
Environment.getJobConfiguration()
Returns the job-wide configuration object that was attached to the JobGraph.
|
Configuration |
Environment.getTaskConfiguration()
Returns the task-wide configuration object, originally attache to the job vertex.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
JobInformation.getJobConfiguration() |
Configuration |
ExecutionGraph.getJobConfiguration() |
Configuration |
TaskInformation.getTaskConfiguration() |
Constructor and Description |
---|
ExecutionGraph(Executor futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
scala.concurrent.duration.FiniteDuration timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
Scheduler scheduler,
ClassLoader userClassLoader,
MetricGroup metricGroup) |
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<BlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
TaskInformation(JobVertexID jobVertexId,
String taskName,
int parallelism,
String invokableClassName,
Configuration taskConfiguration) |
Modifier and Type | Method and Description |
---|---|
static NoRestartStrategy.NoRestartStrategyFactory |
NoRestartStrategy.createFactory(Configuration configuration)
Creates a NoRestartStrategy instance.
|
static FixedDelayRestartStrategy.FixedDelayRestartStrategyFactory |
FixedDelayRestartStrategy.createFactory(Configuration configuration)
Creates a FixedDelayRestartStrategy from the given Configuration.
|
static FailureRateRestartStrategy.FailureRateRestartStrategyFactory |
FailureRateRestartStrategy.createFactory(Configuration configuration) |
static RestartStrategyFactory |
RestartStrategyFactory.createRestartStrategyFactory(Configuration configuration)
Creates a
RestartStrategy instance from the given Configuration . |
Constructor and Description |
---|
FileCache(Configuration config) |
Constructor and Description |
---|
NettyConfig(InetAddress serverAddress,
int serverPort,
int memorySegmentSize,
int numberOfSlots,
Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
JobVertex.getConfiguration()
Returns the vertex's configuration object which can be used to pass custom settings to the task at runtime.
|
Configuration |
JobGraph.getJobConfiguration()
Returns the configuration object for this job.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
AbstractInvokable.getJobConfiguration()
Returns the job configuration object which was attached to the original
JobGraph . |
Configuration |
AbstractInvokable.getTaskConfiguration()
Returns the task configuration object which was attached to the original
JobVertex . |
Modifier and Type | Method and Description |
---|---|
protected Configuration |
JobManager.flinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
static scala.Tuple4<Configuration,JobManagerMode,String,Iterator<Integer>> |
JobManager.parseArgs(String[] args)
Loads the configuration, execution mode and the listening address from the provided command
line arguments.
|
scala.Tuple4<Configuration,JobManagerMode,String,Iterator<Integer>> |
JobManager$.parseArgs(String[] args)
Loads the configuration, execution mode and the listening address from the provided command
line arguments.
|
Modifier and Type | Method and Description |
---|---|
static scala.Tuple12<InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategyFactory,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,SavepointStore,scala.concurrent.duration.FiniteDuration,scala.Option<MetricRegistry>> |
JobManager.createJobManagerComponents(Configuration configuration,
Executor futureExecutor,
Executor ioExecutor,
scala.Option<LeaderElectionService> leaderElectionServiceOption)
Create the job manager components as (instanceManager, scheduler, libraryCacheManager,
archiverProps, defaultExecutionRetries,
delayBetweenRetries, timeout)
|
scala.Tuple12<InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategyFactory,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,SavepointStore,scala.concurrent.duration.FiniteDuration,scala.Option<MetricRegistry>> |
JobManager$.createJobManagerComponents(Configuration configuration,
Executor futureExecutor,
Executor ioExecutor,
scala.Option<LeaderElectionService> leaderElectionServiceOption)
Create the job manager components as (instanceManager, scheduler, libraryCacheManager,
archiverProps, defaultExecutionRetries,
delayBetweenRetries, timeout)
|
static RecoveryMode |
RecoveryMode.fromConfig(Configuration config)
Return the configured
RecoveryMode . |
static akka.actor.ActorRef |
JobManager.getJobManagerActorRef(InetSocketAddress address,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
akka.actor.ActorRef |
JobManager$.getJobManagerActorRef(InetSocketAddress address,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
static String |
JobManager.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
String |
JobManager$.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
static boolean |
RecoveryMode.isHighAvailabilityModeActivated(Configuration configuration)
Returns true if the defined recovery mode supports high availability.
|
static void |
JobManager.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort)
Starts and runs the JobManager with all its components.
|
void |
JobManager$.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort)
Starts and runs the JobManager with all its components.
|
static void |
JobManager.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
Iterator<Integer> listeningPortRange)
Starts and runs the JobManager with all its components trying to bind to
a port in the specified range.
|
void |
JobManager$.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
Iterator<Integer> listeningPortRange)
Starts and runs the JobManager with all its components trying to bind to
a port in the specified range.
|
static scala.Tuple5<akka.actor.ActorSystem,akka.actor.ActorRef,akka.actor.ActorRef,scala.Option<WebMonitor>,scala.Option<akka.actor.ActorRef>> |
JobManager.startActorSystemAndJobManagerActors(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort,
Executor futureExecutor,
Executor ioExecutor,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass,
scala.Option<Class<? extends FlinkResourceManager<?>>> resourceManagerClass)
Starts an ActorSystem, the JobManager and all its components including the WebMonitor.
|
scala.Tuple5<akka.actor.ActorSystem,akka.actor.ActorRef,akka.actor.ActorRef,scala.Option<WebMonitor>,scala.Option<akka.actor.ActorRef>> |
JobManager$.startActorSystemAndJobManagerActors(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort,
Executor futureExecutor,
Executor ioExecutor,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass,
scala.Option<Class<? extends FlinkResourceManager<?>>> resourceManagerClass)
Starts an ActorSystem, the JobManager and all its components including the WebMonitor.
|
static scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Executor futureExecutor,
Executor ioExecutor,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in
the given actor system.
|
scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager$.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Executor futureExecutor,
Executor ioExecutor,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in
the given actor system.
|
static scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Executor futureExecutor,
Executor ioExecutor,
scala.Option<String> jobManagerActorName,
scala.Option<String> archiveActorName,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager$.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Executor futureExecutor,
Executor ioExecutor,
scala.Option<String> jobManagerActorName,
scala.Option<String> archiveActorName,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
Constructor and Description |
---|
JobManager(Configuration flinkConfiguration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
Constructor and Description |
---|
MetricRegistry(Configuration config)
Creates a new MetricRegistry and starts the configured reporter.
|
Modifier and Type | Method and Description |
---|---|
static ScopeFormats |
ScopeFormats.fromConfig(Configuration config)
Creates the scope formats as defined in the given configuration
|
Modifier and Type | Method and Description |
---|---|
Configuration |
FlinkMiniCluster.configuration() |
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Configuration |
LocalFlinkMiniCluster.getDefaultConfig() |
Configuration |
FlinkMiniCluster.userConfiguration() |
Modifier and Type | Method and Description |
---|---|
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
void |
LocalFlinkMiniCluster.initializeIOFormatClasses(Configuration configuration) |
void |
FlinkMiniCluster.setDefaultCiConfig(Configuration config)
Sets CI environment (Travis) specific config defaults.
|
void |
LocalFlinkMiniCluster.setMemory(Configuration config) |
scala.Option<WebMonitor> |
FlinkMiniCluster.startWebServer(Configuration config,
akka.actor.ActorSystem actorSystem,
String jobManagerAkkaURL) |
Constructor and Description |
---|
FlinkMiniCluster(Configuration userConfiguration,
boolean useSingleActorSystem) |
LocalFlinkMiniCluster(Configuration userConfiguration) |
LocalFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
Modifier and Type | Method and Description |
---|---|
static void |
BatchTask.openUserCode(Function stub,
Configuration parameters)
Opens the given stub using its
RichFunction.open(Configuration) method. |
Modifier and Type | Method and Description |
---|---|
void |
CombiningUnilateralSortMerger.setUdfConfiguration(Configuration config) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
TaskConfig.config |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskConfig.getConfiguration()
Gets the configuration that holds the actual values encoded.
|
Configuration |
TaskConfig.getStubParameters() |
Modifier and Type | Method and Description |
---|---|
void |
TaskConfig.setStubParameters(Configuration parameters) |
Constructor and Description |
---|
TaskConfig(Configuration config)
Creates a new Task Config that wraps the given configuration.
|
Modifier and Type | Method and Description |
---|---|
AbstractStateBackend |
StateBackendFactory.createFromConfig(Configuration config)
Creates the state backend, optionally using the given configuration.
|
Modifier and Type | Method and Description |
---|---|
FsStateBackend |
FsStateBackendFactory.createFromConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskManagerConfiguration.configuration() |
Configuration |
TaskManagerRuntimeInfo.getConfiguration()
Gets the configuration that the TaskManager was started with.
|
Configuration |
Task.getJobConfiguration() |
Configuration |
RuntimeEnvironment.getJobConfiguration() |
Configuration |
Task.getTaskConfiguration() |
Configuration |
RuntimeEnvironment.getTaskConfiguration() |
Configuration |
TaskManager$.parseArgsAndLoadConfig(String[] args)
Parse the command line arguments of the TaskManager and loads the configuration.
|
static Configuration |
TaskManager.parseArgsAndLoadConfig(String[] args)
Parse the command line arguments of the TaskManager and loads the configuration.
|
Modifier and Type | Method and Description |
---|---|
scala.Tuple2<String,Object> |
TaskManager$.getAndCheckJobManagerAddress(Configuration configuration)
Gets the hostname and port of the JobManager from the configuration.
|
static scala.Tuple2<String,Object> |
TaskManager.getAndCheckJobManagerAddress(Configuration configuration)
Gets the hostname and port of the JobManager from the configuration.
|
scala.Tuple4<TaskManagerConfiguration,NetworkEnvironmentConfiguration,InstanceConnectionInfo,MemoryType> |
TaskManager$.parseTaskManagerConfiguration(Configuration configuration,
String taskManagerHostname,
boolean localTaskManagerCommunication)
Utility method to extract TaskManager config parameters from the configuration and to
sanity check them.
|
static scala.Tuple4<TaskManagerConfiguration,NetworkEnvironmentConfiguration,InstanceConnectionInfo,MemoryType> |
TaskManager.parseTaskManagerConfiguration(Configuration configuration,
String taskManagerHostname,
boolean localTaskManagerCommunication)
Utility method to extract TaskManager config parameters from the configuration and to
sanity check them.
|
void |
TaskManager$.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
void |
TaskManager$.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
scala.Tuple2<String,Object> |
TaskManager$.selectNetworkInterfaceAndPort(Configuration configuration) |
static scala.Tuple2<String,Object> |
TaskManager.selectNetworkInterfaceAndPort(Configuration configuration) |
void |
TaskManager$.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
ResourceID resourceID,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
ResourceID resourceID,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
akka.actor.ActorRef |
TaskManager$.startTaskManagerComponentsAndActor(Configuration configuration,
ResourceID resourceID,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass) |
static akka.actor.ActorRef |
TaskManager.startTaskManagerComponentsAndActor(Configuration configuration,
ResourceID resourceID,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass) |
Constructor and Description |
---|
RuntimeEnvironment(JobID jobId,
JobVertexID jobVertexId,
ExecutionAttemptID executionId,
ExecutionConfig executionConfig,
TaskInfo taskInfo,
Configuration jobConfiguration,
Configuration taskConfiguration,
ClassLoader userCodeClassLoader,
MemoryManager memManager,
IOManager ioManager,
BroadcastVariableManager bcVarManager,
AccumulatorRegistry accumulatorRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
ActorGateway jobManager,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
TaskManagerConfiguration(String[] tmpDirPaths,
long cleanupInterval,
scala.concurrent.duration.FiniteDuration timeout,
scala.Option<scala.concurrent.duration.FiniteDuration> maxRegistrationDuration,
int numberOfSlots,
Configuration configuration) |
TaskManagerConfiguration(String[] tmpDirPaths,
long cleanupInterval,
scala.concurrent.duration.FiniteDuration timeout,
scala.Option<scala.concurrent.duration.FiniteDuration> maxRegistrationDuration,
int numberOfSlots,
Configuration configuration,
scala.concurrent.duration.FiniteDuration initialRegistrationPause,
scala.concurrent.duration.FiniteDuration maxRegistrationPause,
scala.concurrent.duration.FiniteDuration refusedRegistrationPause) |
TaskManagerRuntimeInfo(String hostname,
Configuration configuration,
String tmpDirectory)
Creates a runtime info.
|
TaskManagerRuntimeInfo(String hostname,
Configuration configuration,
String[] tmpDirectories)
Creates a runtime info.
|
Constructor and Description |
---|
TestingJobManager(Configuration flinkConfiguration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricRegistry) |
Constructor and Description |
---|
TestingResourceManager(Configuration flinkConfig,
LeaderRetrievalService leaderRetriever) |
Modifier and Type | Method and Description |
---|---|
static ZooKeeperCheckpointIDCounter |
ZooKeeperUtils.createCheckpointIDCounter(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId)
Creates a
ZooKeeperCheckpointIDCounter instance. |
static CompletedCheckpointStore |
ZooKeeperUtils.createCompletedCheckpoints(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId,
int maxNumberOfCheckpointsToRetain,
ClassLoader userClassLoader,
Executor executor)
Creates a
ZooKeeperCompletedCheckpointStore instance. |
static ZooKeeperLeaderElectionService |
ZooKeeperUtils.createLeaderElectionService(Configuration configuration)
Creates a
ZooKeeperLeaderElectionService instance and a new CuratorFramework client. |
static ZooKeeperLeaderElectionService |
ZooKeeperUtils.createLeaderElectionService(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperLeaderElectionService instance. |
static ZooKeeperLeaderRetrievalService |
ZooKeeperUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
ZooKeeperLeaderRetrievalService instance. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
StandaloneLeaderRetrievalService from the given configuration. |
static LeaderRetrievalService |
LeaderRetrievalUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
LeaderRetrievalService based on the provided Configuration object. |
static LeaderRetrievalService |
LeaderRetrievalUtils.createLeaderRetrievalService(Configuration configuration,
akka.actor.ActorRef standaloneRef)
Creates a
LeaderRetrievalService that either uses the distributed leader election
configured in the configuration, or, in standalone mode, the given actor reference. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration,
String jobManagerName)
Creates a
StandaloneLeaderRetrievalService form the given configuration and the
JobManager name. |
static ZooKeeperSubmittedJobGraphStore |
ZooKeeperUtils.createSubmittedJobGraphs(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
Executor executor)
Creates a
ZooKeeperSubmittedJobGraphStore instance. |
static RecoveryMode |
LeaderRetrievalUtils.getRecoveryMode(Configuration config)
Gets the recovery mode as configured, based on the
ConfigConstants.RECOVERY_MODE
config key. |
static String |
ZooKeeperUtils.getZooKeeperEnsemble(Configuration flinkConf)
Returns the configured ZooKeeper quorum (and removes whitespace, because ZooKeeper does not
tolerate it).
|
static boolean |
ZooKeeperUtils.isZooKeeperRecoveryMode(Configuration flinkConf)
Returns whether
RecoveryMode.ZOOKEEPER is configured. |
static org.apache.curator.framework.CuratorFramework |
ZooKeeperUtils.startCuratorFramework(Configuration configuration)
Starts a
CuratorFramework instance and connects it to the given ZooKeeper
quorum. |
Modifier and Type | Method and Description |
---|---|
static WebMonitorUtils.LogFileLocation |
WebMonitorUtils.LogFileLocation.find(Configuration config)
Finds the Flink log directory using log.file Java property that is set during startup.
|
static WebMonitor |
WebMonitorUtils.startWebRuntimeMonitor(Configuration config,
LeaderRetrievalService leaderRetrievalService,
akka.actor.ActorSystem actorSystem)
Starts the web runtime monitor.
|
Constructor and Description |
---|
WebMonitorConfig(Configuration config) |
WebRuntimeMonitor(Configuration config,
LeaderRetrievalService leaderRetrievalService,
akka.actor.ActorSystem actorSystem) |
Constructor and Description |
---|
JobManagerConfigHandler(Configuration config) |
TaskManagerLogHandler(JobManagerRetriever retriever,
scala.concurrent.ExecutionContextExecutor executor,
scala.concurrent.Future<String> localJobManagerAddressPromise,
scala.concurrent.duration.FiniteDuration timeout,
TaskManagerLogHandler.FileMode fileMode,
Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
RemoteStreamEnvironment.getClientConfiguration() |
Modifier and Type | Method and Description |
---|---|
static LocalStreamEnvironment |
StreamExecutionEnvironment.createLocalEnvironment(int parallelism,
Configuration configuration)
Creates a
LocalStreamEnvironment . |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfig,
String... jarFiles)
Creates a
RemoteStreamEnvironment . |
Constructor and Description |
---|
LocalStreamEnvironment(Configuration config)
Creates a new local stream environment that configures its local executor with the given configuration.
|
RemoteStreamEnvironment(String host,
int port,
Configuration clientConfiguration,
String... jarFiles)
Creates a new RemoteStreamEnvironment that points to the master
(JobManager) described by the given host name and port.
|
RemoteStreamEnvironment(String host,
int port,
Configuration clientConfiguration,
String[] jarFiles,
URL[] globalClasspaths)
Creates a new RemoteStreamEnvironment that points to the master
(JobManager) described by the given host name and port.
|
Modifier and Type | Method and Description |
---|---|
void |
SocketClientSink.open(Configuration parameters)
Initialize the connection with the Socket in the server.
|
void |
PrintSinkFunction.open(Configuration parameters) |
void |
OutputFormatSinkFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
MultipleIdsMessageAcknowledgingSourceBase.open(Configuration parameters) |
void |
MessageAcknowledgingSourceBase.open(Configuration parameters) |
void |
InputFormatSourceFunction.open(Configuration parameters) |
void |
FromSplittableIteratorFunction.open(Configuration parameters) |
void |
ContinuousFileMonitoringFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
FoldApplyWindowFunction.open(Configuration configuration) |
void |
FoldApplyAllWindowFunction.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
StreamConfig.getConfiguration() |
Constructor and Description |
---|
StreamConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
AbstractUdfStreamOperator.getUserFunctionParameters()
Since the streaming API does not implement any parametrization of functions via a
configuration, the config returned here is actually empty.
|
Modifier and Type | Method and Description |
---|---|
void |
StatefulFunction.open(Configuration c) |
Modifier and Type | Method and Description |
---|---|
void |
ScalaWindowFunctionWrapper.open(Configuration parameters) |
void |
ScalaAllWindowFunctionWrapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
CassandraSinkBase.open(Configuration configuration) |
void |
CassandraPojoSink.open(Configuration configuration) |
void |
CassandraTupleSink.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
void |
ElasticsearchSink.open(Configuration configuration)
Initializes the connection to Elasticsearch by either creating an embedded
Node and retrieving the
Client from it or by creating a
TransportClient . |
Modifier and Type | Method and Description |
---|---|
void |
ElasticsearchSink.open(Configuration configuration)
Initializes the connection to Elasticsearch by creating a
TransportClient . |
Modifier and Type | Method and Description |
---|---|
void |
FlumeSink.open(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
RollingSink.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
FlinkKafkaProducerBase.open(Configuration configuration)
Initializes the connection to Kafka.
|
Modifier and Type | Method and Description |
---|---|
void |
NiFiSink.open(Configuration parameters) |
void |
NiFiSource.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
RMQSource.open(Configuration config) |
void |
RMQSink.open(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
RedisSink.open(Configuration parameters)
Initializes the connection to Redis by either cluster or sentinels or single server.
|
Modifier and Type | Method and Description |
---|---|
void |
TwitterSource.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
WikipediaEditsSource.open(Configuration parameters) |
Constructor and Description |
---|
StreamInputProcessor(InputGate[] inputGates,
TypeSerializer<IN> inputSerializer,
StatefulTask<?> checkpointListener,
CheckpointingMode checkpointMode,
IOManager ioManager,
boolean enableWatermarkMultiplexing,
Configuration taskManagerConfig) |
StreamTwoInputProcessor(Collection<InputGate> inputGates1,
Collection<InputGate> inputGates2,
TypeSerializer<IN1> inputSerializer1,
TypeSerializer<IN2> inputSerializer2,
StatefulTask<?> checkpointListener,
CheckpointingMode checkpointMode,
IOManager ioManager,
boolean enableWatermarkMultiplexing,
Configuration taskManagerConfig) |
Modifier and Type | Method and Description |
---|---|
void |
InternalSingleValueWindowFunction.open(Configuration parameters) |
void |
InternalSingleValueAllWindowFunction.open(Configuration parameters) |
void |
InternalIterableWindowFunction.open(Configuration parameters) |
void |
InternalIterableAllWindowFunction.open(Configuration parameters) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
AbstractTestBase.config
Configuration to start the testing cluster with
|
Modifier and Type | Method and Description |
---|---|
Configuration |
ForkableFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ForkableFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
static ForkableFlinkMiniCluster |
TestBaseUtils.startCluster(Configuration config,
boolean singleActorSystem) |
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(Configuration... testConfigs) |
Modifier and Type | Method and Description |
---|---|
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(List<Configuration> testConfigs) |
Constructor and Description |
---|
AbstractTestBase(Configuration config) |
ForkableFlinkMiniCluster(Configuration userConfiguration) |
ForkableFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
JavaProgramTestBase(Configuration config) |
Modifier and Type | Method and Description |
---|---|
static <T> T |
InstantiationUtil.readObjectFromConfig(Configuration config,
String key,
ClassLoader cl) |
static void |
InstantiationUtil.writeObjectToConfig(Object o,
Configuration config,
String key) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ApplicationClient.flinkConfig() |
Configuration |
YarnClusterClient.getFlinkConfiguration() |
Configuration |
AbstractYarnClusterDescriptor.getFlinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
static int |
Utils.calculateHeapSize(int memory,
Configuration conf)
See documentation
|
static akka.actor.Props |
YarnFlinkResourceManager.createActorProps(Class<? extends YarnFlinkResourceManager> actorClass,
Configuration flinkConfig,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfig,
LeaderRetrievalService leaderRetrievalService,
String applicationMasterHostName,
String webFrontendURL,
ContaineredTaskManagerParameters taskManagerParameters,
org.apache.hadoop.yarn.api.records.ContainerLaunchContext taskManagerLaunchContext,
int numInitialTaskManagers,
org.slf4j.Logger log)
Creates the props needed to instantiate this actor.
|
static org.apache.hadoop.yarn.api.records.ContainerLaunchContext |
YarnApplicationMasterRunner.createTaskManagerContext(Configuration flinkConfig,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfig,
Map<String,String> env,
ContaineredTaskManagerParameters tmParams,
Configuration taskManagerConfig,
String workingDirectory,
Class<?> taskManagerMainClass,
org.slf4j.Logger log)
Creates the launch context, which describes how to bring up a TaskManager process in
an allocated YARN container.
|
protected YarnClusterClient |
AbstractYarnClusterDescriptor.createYarnClusterClient(AbstractYarnClusterDescriptor descriptor,
org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
org.apache.hadoop.yarn.api.records.ApplicationReport report,
Configuration flinkConfiguration,
org.apache.hadoop.fs.Path sessionFilesDir,
boolean perJobCluster)
Creates a YarnClusterClient; may be overriden in tests
|
static Map<String,String> |
Utils.getEnvironmentVariables(String envPrefix,
Configuration flinkConfiguration)
Method to extract environment variables from the flinkConfiguration based on the given prefix String.
|
void |
AbstractYarnClusterDescriptor.setFlinkConfiguration(Configuration conf) |
Constructor and Description |
---|
ApplicationClient(Configuration flinkConfig,
LeaderRetrievalService leaderRetrievalService) |
YarnClusterClient(AbstractYarnClusterDescriptor clusterDescriptor,
org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
org.apache.hadoop.yarn.api.records.ApplicationReport appReport,
Configuration flinkConfig,
org.apache.hadoop.fs.Path sessionFilesDir,
boolean newlyCreatedCluster)
Create a new Flink on YARN cluster.
|
YarnFlinkResourceManager(Configuration flinkConfig,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfig,
LeaderRetrievalService leaderRetrievalService,
String applicationMasterHostName,
String webInterfaceURL,
ContaineredTaskManagerParameters taskManagerParameters,
org.apache.hadoop.yarn.api.records.ContainerLaunchContext taskManagerLaunchContext,
int yarnHeartbeatIntervalMillis,
int maxFailedContainers,
int numInitialTaskManagers) |
YarnFlinkResourceManager(Configuration flinkConfig,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfig,
LeaderRetrievalService leaderRetrievalService,
String applicationMasterHostName,
String webInterfaceURL,
ContaineredTaskManagerParameters taskManagerParameters,
org.apache.hadoop.yarn.api.records.ContainerLaunchContext taskManagerLaunchContext,
int yarnHeartbeatIntervalMillis,
int maxFailedContainers,
int numInitialTaskManagers,
YarnResourceManagerCallbackHandler callbackHandler) |
YarnFlinkResourceManager(Configuration flinkConfig,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfig,
LeaderRetrievalService leaderRetrievalService,
String applicationMasterHostName,
String webInterfaceURL,
ContaineredTaskManagerParameters taskManagerParameters,
org.apache.hadoop.yarn.api.records.ContainerLaunchContext taskManagerLaunchContext,
int yarnHeartbeatIntervalMillis,
int maxFailedContainers,
int numInitialTaskManagers,
YarnResourceManagerCallbackHandler callbackHandler,
org.apache.hadoop.yarn.client.api.async.AMRMClientAsync<org.apache.hadoop.yarn.client.api.AMRMClient.ContainerRequest> resourceManagerClient,
org.apache.hadoop.yarn.client.api.NMClient nodeManagerClient) |
YarnJobManager(Configuration flinkConfiguration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
Modifier and Type | Method and Description |
---|---|
YarnClusterClient |
FlinkYarnSessionCli.createCluster(String applicationName,
org.apache.commons.cli.CommandLine cmdLine,
Configuration config) |
static File |
FlinkYarnSessionCli.getYarnPropertiesLocation(Configuration conf) |
boolean |
FlinkYarnSessionCli.isActive(org.apache.commons.cli.CommandLine commandLine,
Configuration configuration) |
YarnClusterClient |
FlinkYarnSessionCli.retrieveCluster(org.apache.commons.cli.CommandLine cmdLine,
Configuration config) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.