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 and the string that
defines the record delimiter.
|
void |
BinaryOutputFormat.configure(Configuration parameters) |
void |
BinaryInputFormat.configure(Configuration parameters) |
protected static void |
DelimitedInputFormat.loadConfigParameters(Configuration parameters) |
Constructor and Description |
---|
DelimitedInputFormat(Path filePath,
Configuration configuration) |
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.createLocalEnvironmentWithWebUI(Configuration conf)
Creates a
LocalEnvironment for local program execution that also starts the
web monitoring UI. |
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 |
---|---|
ExecutionEnvironment |
ExecutionEnvironment$.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
static ExecutionEnvironment |
ExecutionEnvironment.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
ExecutionEnvironment |
ExecutionEnvironment$.createLocalEnvironmentWithWebUI(Configuration config)
Creates a
ExecutionEnvironment for local program execution that also starts the
web monitoring UI. |
static ExecutionEnvironment |
ExecutionEnvironment.createLocalEnvironmentWithWebUI(Configuration config)
Creates a
ExecutionEnvironment for local program execution that also starts the
web monitoring UI. |
ExecutionEnvironment |
ExecutionEnvironment$.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
scala.collection.Seq<String> jarFiles)
Creates a remote execution environment.
|
static 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 |
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,
List<URL> userJarFiles) |
ClusterType |
CustomCommandLine.createCluster(String applicationName,
org.apache.commons.cli.CommandLine commandLine,
Configuration config,
List<URL> userJarFiles)
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 |
DelegatingConfiguration.clone() |
Configuration |
Configuration.clone() |
static Configuration |
GlobalConfiguration.getDynamicProperties()
Get the dynamic properties.
|
static Configuration |
GlobalConfiguration.loadConfiguration()
Loads the global configuration from the environment.
|
static Configuration |
GlobalConfiguration.loadConfiguration(String configDir)
Loads the configuration files from the specified directory.
|
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.setDynamicProperties(Configuration dynamicProperties)
Set the process-wide dynamic properties to be merged with the loaded 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 |
---|---|
AbstractStateBackend |
RocksDBStateBackendFactory.createFromConfig(Configuration config) |
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 |
AnalyticHelper.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopReduceFunction.open(Configuration parameters) |
void |
HadoopReduceCombineFunction.open(Configuration parameters) |
void |
HadoopMapFunction.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) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
FlinkMesosSessionCli.decodeDynamicProperties(String dynamicPropertiesEncoded)
Decode encoded dynamic properties.
|
Modifier and Type | Method and Description |
---|---|
static String |
FlinkMesosSessionCli.encodeDynamicProperties(Configuration configuration)
Encode dynamic properties as a string to be transported as an environment variable.
|
Modifier and Type | Method and Description |
---|---|
static MesosTaskManagerParameters |
MesosTaskManagerParameters.create(Configuration flinkConfig)
Create the Mesos TaskManager parameters.
|
static akka.actor.Props |
MesosFlinkResourceManager.createActorProps(Class<? extends MesosFlinkResourceManager> actorClass,
Configuration flinkConfig,
MesosConfiguration mesosConfig,
MesosWorkerStore workerStore,
LeaderRetrievalService leaderRetrievalService,
MesosTaskManagerParameters taskManagerParameters,
ContainerSpecification taskManagerContainerSpec,
MesosArtifactResolver artifactResolver,
org.slf4j.Logger log)
Creates the props needed to instantiate this actor.
|
static MesosConfiguration |
MesosApplicationMasterRunner.createMesosConfig(Configuration flinkConfig,
String hostname)
Loads and validates the ResourceManager Mesos configuration from the given Flink configuration.
|
protected int |
MesosApplicationMasterRunner.runPrivileged(Configuration config,
Configuration dynamicProperties)
The main work method, must run as a privileged action.
|
Constructor and Description |
---|
MesosFlinkResourceManager(Configuration flinkConfig,
MesosConfiguration mesosConfig,
MesosWorkerStore workerStore,
LeaderRetrievalService leaderRetrievalService,
MesosTaskManagerParameters taskManagerParameters,
ContainerSpecification taskManagerContainerSpec,
MesosArtifactResolver artifactResolver,
int maxFailedTasks,
int numInitialTaskManagers) |
MesosJobManager(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,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
Modifier and Type | Method and Description |
---|---|
static MesosServices |
MesosServicesUtils.createMesosServices(Configuration configuration)
Creates a
MesosServices instance depending on the high availability settings. |
MesosWorkerStore |
ZooKeeperMesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor) |
MesosWorkerStore |
StandaloneMesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor) |
MesosWorkerStore |
MesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor)
Creates a
MesosWorkerStore which is used to persist mesos worker in high availability
mode. |
Modifier and Type | Method and Description |
---|---|
<T extends LaunchCoordinator> |
LaunchCoordinator$.createActorProps(Class<T> actorClass,
akka.actor.ActorRef manager,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskSchedulerBuilder optimizerBuilder)
Get the configuration properties for the launch coordinator.
|
static <T extends LaunchCoordinator> |
LaunchCoordinator.createActorProps(Class<T> actorClass,
akka.actor.ActorRef manager,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskSchedulerBuilder optimizerBuilder)
Get the configuration properties for the launch coordinator.
|
static <T extends ConnectionMonitor> |
ConnectionMonitor.createActorProps(Class<T> actorClass,
Configuration flinkConfig)
Creates the properties for the ConnectionMonitor actor.
|
<T extends ConnectionMonitor> |
ConnectionMonitor$.createActorProps(Class<T> actorClass,
Configuration flinkConfig)
Creates the properties for the ConnectionMonitor actor.
|
static <T extends ReconciliationCoordinator> |
ReconciliationCoordinator.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver)
Create the properties for a reconciliation coordinator.
|
<T extends ReconciliationCoordinator> |
ReconciliationCoordinator$.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver)
Create the properties for a reconciliation coordinator.
|
<T extends Tasks,M extends TaskMonitor> |
Tasks$.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
Class<M> taskMonitorClass)
Create a tasks actor.
|
static <T extends Tasks,M extends TaskMonitor> |
Tasks.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
Class<M> taskMonitorClass)
Create a tasks actor.
|
<T extends TaskMonitor> |
TaskMonitor$.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskMonitor.TaskGoalState goalState)
Creates the properties for the TaskMonitor actor.
|
static <T extends TaskMonitor> |
TaskMonitor.createActorProps(Class<T> actorClass,
Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskMonitor.TaskGoalState goalState)
Creates the properties for the TaskMonitor actor.
|
Constructor and Description |
---|
LaunchCoordinator(akka.actor.ActorRef manager,
Configuration config,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskSchedulerBuilder optimizerBuilder) |
ReconciliationCoordinator(Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver) |
TaskMonitor(Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
TaskMonitor.TaskGoalState goalState) |
Tasks(Configuration flinkConfig,
org.apache.mesos.SchedulerDriver schedulerDriver,
scala.Function2<akka.actor.ActorRefFactory,TaskMonitor.TaskGoalState,akka.actor.ActorRef> taskMonitorCreator) |
Modifier and Type | Method and Description |
---|---|
static org.apache.curator.framework.CuratorFramework |
ZooKeeperUtils.startCuratorFramework(Configuration configuration)
Starts a
CuratorFramework instance and connects it to the given ZooKeeper
quorum. |
Constructor and Description |
---|
MesosArtifactServer(String prefix,
String serverHostname,
int configuredPort,
Configuration config) |
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>> externalAddress)
Creates an akka config with the provided configuration values.
|
static com.typesafe.config.Config |
AkkaUtils.getAkkaConfig(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> externalAddress)
Creates an akka config with the provided configuration values.
|
String |
AkkaUtils$.getAkkaProtocol(Configuration config)
Returns the protocol field for the URL of the remote actor system given the user configuration
|
static String |
AkkaUtils.getAkkaProtocol(Configuration config)
Returns the protocol field for the URL of the remote actor system given the user configuration
|
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) |
Modifier and Type | Method and Description |
---|---|
static List<BlobKey> |
BlobClient.uploadJarFiles(ActorGateway jobManager,
scala.concurrent.duration.FiniteDuration askTimeout,
Configuration clientConfig,
List<Path> jars)
Retrieves the
BlobServer address from the JobManager and uploads
the JAR files to it. |
static List<BlobKey> |
BlobClient.uploadJarFiles(InetSocketAddress serverAddress,
Configuration clientConfig,
List<Path> jars)
Uploads the JAR files to a
BlobServer at the given address. |
Constructor and Description |
---|
BlobCache(InetSocketAddress serverAddress,
Configuration blobClientConfig) |
BlobClient(InetSocketAddress serverAddress,
Configuration clientConfig)
Instantiates a new BLOB client.
|
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 JobListeningContext |
JobClient.attachToRunningJob(JobID jobID,
ActorGateway jobManagerGateWay,
Configuration configuration,
akka.actor.ActorSystem actorSystem,
LeaderRetrievalService leaderRetrievalService,
scala.concurrent.duration.FiniteDuration timeout,
boolean sysoutLogUpdates)
Attaches to a running Job using the JobID.
|
static akka.actor.Props |
JobSubmissionClientActor.createActorProps(LeaderRetrievalService leaderRetrievalService,
scala.concurrent.duration.FiniteDuration timeout,
boolean sysoutUpdates,
Configuration clientConfig) |
static ClassLoader |
JobClient.retrieveClassLoader(JobID jobID,
ActorGateway jobManager,
Configuration config)
Reconstructs the class loader by first requesting information about it at the JobManager
and then downloading missing jar files.
|
static akka.actor.ActorSystem |
JobClient.startJobClientActorSystem(Configuration config) |
static JobListeningContext |
JobClient.submitJob(akka.actor.ActorSystem actorSystem,
Configuration config,
LeaderRetrievalService leaderRetrievalService,
JobGraph jobGraph,
scala.concurrent.duration.FiniteDuration timeout,
boolean sysoutLogUpdates,
ClassLoader classLoader)
Submits a job to a Flink cluster (non-blocking) and returns a JobListeningContext which can be
passed to
awaitJobResult to get the result of the submission. |
static JobExecutionResult |
JobClient.submitJobAndWait(akka.actor.ActorSystem actorSystem,
Configuration config,
LeaderRetrievalService leaderRetrievalService,
JobGraph jobGraph,
scala.concurrent.duration.FiniteDuration timeout,
boolean sysoutLogUpdates,
ClassLoader classLoader)
Sends a [[JobGraph]] to the JobClient actor specified by jobClient which submits it then to
the JobManager.
|
static void |
JobClient.submitJobDetached(ActorGateway jobManagerGateway,
Configuration config,
JobGraph jobGraph,
scala.concurrent.duration.FiniteDuration timeout,
ClassLoader classLoader)
Submits a job in detached mode.
|
Constructor and Description |
---|
JobListeningContext(JobID jobID,
scala.concurrent.Future<Object> jobResultFuture,
akka.actor.ActorRef jobClientActor,
scala.concurrent.duration.FiniteDuration timeout,
akka.actor.ActorSystem actorSystem,
Configuration configuration)
Constructor to use when the class loader is not available.
|
JobSubmissionClientActor(LeaderRetrievalService leaderRetrievalService,
scala.concurrent.duration.FiniteDuration timeout,
boolean sysoutUpdates,
Configuration clientConfig) |
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.
|
Configuration |
ContainerSpecification.getDynamicConfiguration()
Get the dynamic configuration.
|
Configuration |
ContainerSpecification.getSystemProperties()
Get the system properties.
|
static Configuration |
BootstrapTools.parseDynamicProperties(org.apache.commons.cli.CommandLine cmd)
Parse the dynamic properties (passed on the command line).
|
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 String |
ContainerSpecification.formatSystemProperties(Configuration jvmArgs)
Format the system properties as a shell-compatible command-line argument.
|
static Configuration |
BootstrapTools.generateTaskManagerConfiguration(Configuration baseConfig,
String jobManagerHostname,
int jobManagerPort,
int numSlots,
scala.concurrent.duration.FiniteDuration registrationTimeout)
Generate a task manager configuration.
|
static akka.actor.Props |
FlinkResourceManager.getResourceManagerProps(Class<? extends FlinkResourceManager> resourceManagerClass,
Configuration configuration,
LeaderRetrievalService leaderRetrievalService) |
static String |
BootstrapTools.getTaskManagerShellCommand(Configuration flinkConfig,
ContaineredTaskManagerParameters tmParams,
String configDirectory,
String logDirectory,
boolean hasLogback,
boolean hasLog4j,
boolean hasKrb5,
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 |
---|
ContaineredJobManager(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,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
FlinkResourceManager(int numInitialTaskManagers,
Configuration flinkConfig,
LeaderRetrievalService leaderRetriever)
Creates a AbstractFrameworkMaster actor.
|
Modifier and Type | Method and Description |
---|---|
SSLStoreOverlay.Builder |
SSLStoreOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment (and global configuration).
|
Krb5ConfOverlay.Builder |
Krb5ConfOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment.
|
KeytabOverlay.Builder |
KeytabOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment (and global configuration).
|
HadoopUserOverlay.Builder |
HadoopUserOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current Hadoop user information (from
UserGroupInformation ). |
HadoopConfOverlay.Builder |
HadoopConfOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment's Hadoop configuration.
|
FlinkDistributionOverlay.Builder |
FlinkDistributionOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment.
|
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 attached to the job vertex.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
JobInformation.getJobConfiguration() |
Configuration |
ExecutionGraph.getJobConfiguration() |
Configuration |
TaskInformation.getTaskConfiguration() |
Modifier and Type | Method and Description |
---|---|
static ExecutionGraph |
ExecutionGraphBuilder.buildGraph(ExecutionGraph prior,
JobGraph jobGraph,
Configuration jobManagerConfig,
Executor futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time timeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
int parallelismForAutoMax,
org.slf4j.Logger log)
Builds the ExecutionGraph from the JobGraph.
|
Constructor and Description |
---|
ExecutionGraph(Executor futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
Time timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
SlotProvider slotProvider,
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 numberOfSubtasks,
int numberOfKeyGroups,
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 |
---|---|
void |
JobGraph.uploadRequiredJarFiles(InetSocketAddress serverAddress,
Configuration blobClientConfig)
Uploads the previously added user jar file to the job manager through the job manager's BLOB server.
|
void |
JobGraph.uploadUserJars(ActorGateway jobManager,
scala.concurrent.duration.FiniteDuration askTimeout,
Configuration blobClientConfig)
Uploads the previously added user JAR files to the job manager through
the job manager's BLOB server.
|
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 |
---|---|
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.
|
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.
|
Modifier and Type | Method and Description |
---|---|
scala.Tuple11<InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategyFactory,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,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 scala.Tuple11<InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategyFactory,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,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 HighAvailabilityMode |
HighAvailabilityMode.fromConfig(Configuration config)
Return the configured
HighAvailabilityMode . |
akka.actor.ActorRef |
JobManager$.getJobManagerActorRef(String hostPort,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
static akka.actor.ActorRef |
JobManager.getJobManagerActorRef(String hostPort,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
akka.actor.Props |
JobManager$.getJobManagerProps(Class<? extends JobManager> jobManagerClass,
Configuration configuration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphStore,
CheckpointRecoveryFactory checkpointRecoveryFactory,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
static akka.actor.Props |
JobManager.getJobManagerProps(Class<? extends JobManager> jobManagerClass,
Configuration configuration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphStore,
CheckpointRecoveryFactory checkpointRecoveryFactory,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
String |
JobManager$.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
static String |
JobManager.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
static boolean |
HighAvailabilityMode.isHighAvailabilityModeActivated(Configuration configuration)
Returns true if the defined recovery mode supports high availability.
|
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,
int listeningPort)
Starts and runs the JobManager with all its components.
|
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 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.
|
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 externalHostname,
int port,
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.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 externalHostname,
int port,
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.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,
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.
|
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.
|
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,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,Configuration>> |
MetricRegistryConfiguration.getReporterConfigurations() |
Modifier and Type | Method and Description |
---|---|
static MetricRegistryConfiguration |
MetricRegistryConfiguration.fromConfiguration(Configuration configuration)
Create a metric registry configuration object from the given
Configuration . |
Constructor and Description |
---|
MetricRegistryConfiguration(ScopeFormats scopeFormats,
char delimiter,
List<Tuple2<String,Configuration>> reporterConfigurations) |
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() |
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Configuration |
LocalFlinkMiniCluster.getDefaultConfig() |
protected Configuration |
FlinkMiniCluster.originalConfiguration() |
Configuration |
FlinkMiniCluster.userConfiguration() |
Modifier and Type | Method and Description |
---|---|
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
akka.actor.Props |
LocalFlinkMiniCluster.getJobManagerProps(Class<? extends JobManager> jobManagerClass,
Configuration configuration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphStore,
CheckpointRecoveryFactory checkpointRecoveryFactory,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
akka.actor.Props |
LocalFlinkMiniCluster.getResourceManagerProps(Class<? extends FlinkResourceManager<? extends ResourceIDRetrievable>> resourceManagerClass,
Configuration configuration,
LeaderRetrievalService leaderRetrievalService) |
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 SSLContext |
SSLUtils.createSSLClientContext(Configuration sslConfig)
Creates the SSL Context for the client if SSL is configured
|
static SSLContext |
SSLUtils.createSSLServerContext(Configuration sslConfig)
Creates the SSL Context for the server if SSL is configured
|
static boolean |
SSLUtils.getSSLEnabled(Configuration sslConfig)
Retrieves the global ssl flag from configuration
|
static void |
SSLUtils.setSSLVerifyHostname(Configuration sslConfig,
SSLParameters sslParams)
Sets SSL options to verify peer's hostname in the certificate
|
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.
|
Constructor and Description |
---|
QueryableStateClient(Configuration config)
Creates a client from the given configuration.
|
Constructor and Description |
---|
SecurityConfiguration(Configuration flinkConf)
Create a security configuration from the global configuration.
|
SecurityConfiguration(Configuration flinkConf,
org.apache.hadoop.conf.Configuration hadoopConf)
Create a security configuration from the global configuration.
|
SecurityConfiguration(Configuration flinkConf,
org.apache.hadoop.conf.Configuration hadoopConf,
List<? extends Class<? extends SecurityModule>> securityModules)
Create a security configuration from the global 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.Tuple7<TaskManagerConfiguration,TaskManagerLocation,MemoryManager,IOManager,NetworkEnvironment,LeaderRetrievalService,MetricRegistry> |
TaskManager$.createTaskManagerComponents(Configuration configuration,
ResourceID resourceID,
String taskManagerHostname,
boolean localTaskManagerCommunication,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption) |
static scala.Tuple7<TaskManagerConfiguration,TaskManagerLocation,MemoryManager,IOManager,NetworkEnvironment,LeaderRetrievalService,MetricRegistry> |
TaskManager.createTaskManagerComponents(Configuration configuration,
ResourceID resourceID,
String taskManagerHostname,
boolean localTaskManagerCommunication,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption) |
scala.Tuple3<String,String,Object> |
TaskManager$.getAndCheckJobManagerAddress(Configuration configuration)
Gets the protocol, hostname and port of the JobManager from the configuration.
|
static scala.Tuple3<String,String,Object> |
TaskManager.getAndCheckJobManagerAddress(Configuration configuration)
Gets the protocol, hostname and port of the JobManager from the configuration.
|
scala.Tuple4<TaskManagerConfiguration,NetworkEnvironmentConfiguration,InetSocketAddress,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,InetSocketAddress,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)
Starts the task manager actor.
|
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)
Starts the task manager actor.
|
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,
TaskKvStateRegistry kvStateRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
CheckpointResponder checkpointResponder,
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.
|
TaskManagerRuntimeInfo(String hostname,
Configuration configuration,
String[] tmpDirectories,
boolean exitJvmOnOutOfMemory)
Creates a runtime info.
|
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,
Executor executor)
Creates a
ZooKeeperCompletedCheckpointStore instance. |
static <T extends Serializable> |
ZooKeeperUtils.createFileSystemStateStorage(Configuration configuration,
String prefix)
Creates a
FileSystemStateStorageHelper 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,
boolean resolveInitialHostName)
Creates a
StandaloneLeaderRetrievalService from the given configuration. |
static LeaderRetrievalService |
LeaderRetrievalUtils.createLeaderRetrievalService(Configuration configuration,
boolean resolveInitialHostName)
Creates a
LeaderRetrievalService based on the provided Configuration object. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration,
boolean resolveInitialHostName,
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 ZooKeeperUtils.ZkClientACLMode |
ZooKeeperUtils.ZkClientACLMode.fromConfig(Configuration config)
Return the configured
ZooKeeperUtils.ZkClientACLMode . |
static HighAvailabilityMode |
LeaderRetrievalUtils.getRecoveryMode(Configuration config)
Gets the recovery mode as configured, based on the
ConfigConstants.HA_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
HighAvailabilityMode.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 |
---|
JarRunHandler(File jarDirectory,
scala.concurrent.duration.FiniteDuration timeout,
Configuration clientConfig) |
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,
boolean httpsEnabled) |
Constructor and Description |
---|
ZooKeeperUtilityFactory(Configuration configuration,
String path) |
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.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 . |
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 |
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 |
---|---|
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(Configuration config)
Creates a
StreamExecutionEnvironment for local program execution that also starts the
web monitoring UI. |
StreamExecutionEnvironment |
StreamExecutionEnvironment$.createLocalEnvironmentWithWebUI(Configuration config)
Creates a
StreamExecutionEnvironment for local program execution that also starts the
web monitoring UI. |
Modifier and Type | Method and Description |
---|---|
void |
StatefulFunction.open(Configuration c) |
Modifier and Type | Method and Description |
---|---|
void |
CassandraTupleSink.open(Configuration configuration) |
void |
CassandraPojoSink.open(Configuration configuration) |
void |
CassandraSinkBase.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 |
RollingSink.open(Configuration parameters)
Deprecated.
|
RollingSink<T> |
RollingSink.setFSConfig(Configuration config)
Deprecated.
Specify a custom
Configuration that will be used when creating
the FileSystem for writing. |
Modifier and Type | Method and Description |
---|---|
void |
BucketingSink.open(Configuration parameters) |
BucketingSink<T> |
BucketingSink.setFSConfig(Configuration config)
Specify a custom
Configuration that will be used when creating
the FileSystem for writing. |
Modifier and Type | Method and Description |
---|---|
void |
FlinkKafkaProducer010.open(Configuration parameters)
This method is used for approach (a) (see above)
|
void |
FlinkKafkaProducerBase.open(Configuration configuration)
Initializes the connection to Kafka.
|
void |
FlinkKafkaConsumerBase.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
void |
NiFiSource.open(Configuration parameters) |
void |
NiFiSink.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 |
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 checkpointedTask,
CheckpointingMode checkpointMode,
IOManager ioManager,
Configuration taskManagerConfig) |
StreamTwoInputProcessor(Collection<InputGate> inputGates1,
Collection<InputGate> inputGates2,
TypeSerializer<IN1> inputSerializer1,
TypeSerializer<IN2> inputSerializer2,
StatefulTask checkpointedTask,
CheckpointingMode checkpointMode,
IOManager ioManager,
Configuration taskManagerConfig) |
Modifier and Type | Method and Description |
---|---|
void |
CorrelateFlatMapRunner.open(Configuration parameters) |
void |
MapRunner.open(Configuration parameters) |
void |
FlatJoinRunner.open(Configuration parameters) |
void |
LimitFilterFunction.open(Configuration config) |
void |
FlatMapRunner.open(Configuration parameters) |
void |
MapSideJoinRunner.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
IncrementalAggregateAllWindowFunction.open(Configuration parameters) |
void |
IncrementalAggregateWindowFunction.open(Configuration parameters) |
void |
AggregateMapFunction.open(Configuration config) |
void |
IncrementalAggregateTimeWindowFunction.open(Configuration parameters) |
void |
AggregateWindowFunction.open(Configuration parameters) |
void |
AggregateTimeWindowFunction.open(Configuration parameters) |
void |
AggregateAllWindowFunction.open(Configuration parameters) |
void |
AggregateAllTimeWindowFunction.open(Configuration parameters) |
void |
IncrementalAggregateAllTimeWindowFunction.open(Configuration parameters) |
void |
AggregateReduceGroupFunction.open(Configuration config) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
AbstractTestBase.config
Configuration to start the testing cluster with
|
Modifier and Type | Method and Description |
---|---|
static Configuration |
SecureTestEnvironment.populateFlinkSecureConfigurations(Configuration flinkConf) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
SecureTestEnvironment.populateFlinkSecureConfigurations(Configuration flinkConf) |
static LocalFlinkMiniCluster |
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) |
JavaProgramTestBase(Configuration config) |
Modifier and Type | Method and Description |
---|---|
static int |
ConfigurationUtil.getIntegerWithDeprecatedKeys(Configuration config,
String key,
int defaultValue,
String... deprecatedKeys)
Returns the value associated with the given key as an Integer in the
given Configuration.
|
static String |
ConfigurationUtil.getStringWithDeprecatedKeys(Configuration config,
String key,
String defaultValue,
String... deprecatedKeys)
Returns the value associated with the given key as a String in the
given Configuration.
|
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.
|
protected int |
YarnApplicationMasterRunner.runApplicationMaster(Configuration config)
The main work method, must run as a privileged action.
|
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,
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,
List<URL> userJarFiles) |
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.