Modifier and Type | Method and Description |
---|---|
abstract void |
AbstractTableInputFormat.configure(Configuration parameters)
Creates a
Scan object and opens the HTable connection. |
void |
HBaseRowInputFormat.configure(Configuration parameters) |
void |
TableInputFormat.configure(Configuration parameters)
Creates a
Scan object and opens the HTable connection. |
void |
HBaseUpsertSinkFunction.open(Configuration parameters) |
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 |
AbstractRichFunction.open(Configuration parameters) |
void |
RichFunction.open(Configuration parameters)
Initialization method for the function.
|
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 |
GenericInputFormat.configure(Configuration parameters) |
void |
FileInputFormat.configure(Configuration parameters)
Configures the file input format by reading the file path from the configuration.
|
void |
FileOutputFormat.configure(Configuration parameters) |
void |
InputFormat.configure(Configuration parameters)
Configures this input format.
|
void |
OutputFormat.configure(Configuration parameters)
Configures this output format.
|
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) |
static void |
FileOutputFormat.initDefaultsFromConfiguration(Configuration configuration)
Initialize defaults for output format.
|
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,
org.apache.flink.api.scala.FlinkILoop flinkILoop,
Configuration clientConfig,
String... jarFiles)
Creates new ScalaShellRemoteEnvironment that has a reference to the FlinkILoop.
|
ScalaShellRemoteStreamEnvironment(String host,
int port,
org.apache.flink.api.scala.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 |
---|---|
static Configuration |
HadoopUtils.getHadoopConfiguration(Configuration flinkConfiguration)
Returns a new Hadoop Configuration object using the path to the hadoop conf configured
in the main configuration (flink-conf.yaml).
|
Modifier and Type | Method and Description |
---|---|
void |
HadoopOutputFormatBase.configure(Configuration parameters) |
void |
HadoopInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
PrintingOutputFormat.configure(Configuration parameters) |
void |
TextInputFormat.configure(Configuration parameters) |
void |
TextValueInputFormat.configure(Configuration parameters) |
void |
LocalCollectionOutputFormat.configure(Configuration parameters) |
void |
BlockingShuffleOutputFormat.configure(Configuration parameters) |
void |
DiscardingOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractJDBCOutputFormat.configure(Configuration parameters) |
void |
JDBCInputFormat.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 |
SingleInputUdfOperator.getParameters() |
Configuration |
DataSource.getParameters() |
Configuration |
TwoInputUdfOperator.getParameters() |
Configuration |
DataSink.getParameters() |
Modifier and Type | Method and Description |
---|---|
O |
UdfOperator.withParameters(Configuration parameters)
Sets the configuration parameters for the UDF.
|
O |
SingleInputUdfOperator.withParameters(Configuration parameters) |
DataSource<OUT> |
DataSource.withParameters(Configuration parameters)
Pass a configuration to the InputFormat.
|
O |
TwoInputUdfOperator.withParameters(Configuration parameters) |
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 |
RuntimeComparatorFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeSerializerFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeComparatorFactory.writeParametersToConfig(Configuration config) |
void |
RuntimeSerializerFactory.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 |
---|---|
void |
CassandraPojoOutputFormat.configure(Configuration parameters) |
void |
CassandraOutputFormatBase.configure(Configuration parameters) |
void |
CassandraInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
Generator.configure(Configuration parameters) |
void |
BlockingIncrementingMapFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
PatternTimeoutSelectAdapter.open(Configuration parameters) |
void |
PatternSelectAdapter.open(Configuration parameters) |
void |
PatternTimeoutFlatSelectAdapter.open(Configuration parameters) |
void |
PatternFlatSelectAdapter.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
NFA.open(RuntimeContext cepRuntimeContext,
Configuration conf)
Initialization method for the NFA.
|
Modifier and Type | Method and Description |
---|---|
void |
RichCompositeIterativeCondition.open(Configuration parameters) |
void |
RichIterativeCondition.open(Configuration parameters) |
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 | Field and Description |
---|---|
protected Configuration |
AbstractCustomCommandLine.configuration |
Modifier and Type | Method and Description |
---|---|
protected Configuration |
AbstractCustomCommandLine.applyCommandLineOptionsToConfiguration(org.apache.commons.cli.CommandLine commandLine)
Override configuration settings by specified command line options.
|
Configuration |
CliFrontend.getConfiguration()
Getter which returns a copy of the associated configuration.
|
Configuration |
AbstractCustomCommandLine.getConfiguration() |
Modifier and Type | Method and Description |
---|---|
static List<CustomCommandLine<?>> |
CliFrontend.loadCustomCommandLines(Configuration configuration,
String configurationDirectory) |
Constructor and Description |
---|
AbstractCustomCommandLine(Configuration configuration) |
CliFrontend(Configuration configuration,
List<CustomCommandLine<?>> customCommandLines) |
DefaultCLI(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static ClusterSpecification |
ClusterSpecification.fromConfiguration(Configuration configuration) |
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.
|
Modifier and Type | Method and Description |
---|---|
static ClassLoader |
JobWithJars.buildUserCodeClassLoader(List<URL> jars,
List<URL> classpaths,
ClassLoader parent,
Configuration configuration) |
static JobGraph |
PackagedProgramUtils.createJobGraph(PackagedProgram packagedProgram,
Configuration configuration,
int defaultParallelism)
|
static JobGraph |
PackagedProgramUtils.createJobGraph(PackagedProgram packagedProgram,
Configuration configuration,
int defaultParallelism,
JobID jobID)
|
static JobGraph |
ClusterClient.getJobGraph(Configuration flinkConfig,
FlinkPlan optPlan,
List<URL> jarFiles,
List<URL> classpaths,
SavepointRestoreSettings savepointSettings) |
static JobGraph |
ClusterClient.getJobGraph(Configuration flinkConfig,
PackagedProgram prog,
FlinkPlan optPlan,
SavepointRestoreSettings savepointSettings) |
Constructor and Description |
---|
ClusterClient(Configuration flinkConfig)
Creates a instance that submits the programs to the JobManager defined in the
configuration.
|
ClusterClient(Configuration flinkConfig,
HighAvailabilityServices highAvailabilityServices,
boolean sharedHaServices)
Creates a instance that submits the programs to the JobManager defined in the
configuration.
|
MiniClusterClient(Configuration configuration,
MiniCluster miniCluster) |
PackagedProgram(File jarFile,
List<URL> classpaths,
String entryPointClassName,
Configuration configuration,
String... args)
Creates an instance that wraps the plan defined in the jar file using the given
arguments.
|
Modifier and Type | Method and Description |
---|---|
static RestClusterClientConfiguration |
RestClusterClientConfiguration.fromConfiguration(Configuration config) |
Constructor and Description |
---|
RestClusterClient(Configuration config,
T clusterId) |
RestClusterClient(Configuration config,
T clusterId,
LeaderRetrievalService webMonitorRetrievalService) |
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 |
ConfigurationUtils.createConfiguration(Properties properties)
Creates a new
Configuration from the given Properties . |
static Configuration |
GlobalConfiguration.loadConfiguration()
Loads the global configuration from the environment.
|
static Configuration |
GlobalConfiguration.loadConfiguration(Configuration dynamicProperties)
Loads the global configuration and adds the given dynamic properties
configuration.
|
static Configuration |
GlobalConfiguration.loadConfiguration(String configDir)
Loads the configuration files from the specified directory.
|
static Configuration |
GlobalConfiguration.loadConfiguration(String configDir,
Configuration dynamicProperties)
Loads the configuration files from the specified directory.
|
Modifier and Type | Method and Description |
---|---|
void |
DelegatingConfiguration.addAll(Configuration other) |
void |
UnmodifiableConfiguration.addAll(Configuration other) |
void |
Configuration.addAll(Configuration other) |
void |
DelegatingConfiguration.addAll(Configuration other,
String prefix) |
void |
UnmodifiableConfiguration.addAll(Configuration other,
String prefix) |
void |
Configuration.addAll(Configuration other,
String prefix)
Adds all entries from the given configuration into this configuration.
|
static MemorySize |
ConfigurationUtils.getJobManagerHeapMemory(Configuration configuration)
Get job manager's heap memory.
|
static String[] |
CoreOptions.getParentFirstLoaderPatterns(Configuration config) |
static Time |
ConfigurationUtils.getStandaloneClusterStartupPeriodTime(Configuration configuration) |
static Optional<Time> |
ConfigurationUtils.getSystemResourceMetricsProbingInterval(Configuration configuration) |
static MemorySize |
ConfigurationUtils.getTaskManagerHeapMemory(Configuration configuration)
Get task manager's heap memory.
|
static Configuration |
GlobalConfiguration.loadConfiguration(Configuration dynamicProperties)
Loads the global configuration and adds the given dynamic properties
configuration.
|
static Configuration |
GlobalConfiguration.loadConfiguration(String configDir,
Configuration dynamicProperties)
Loads the configuration files from the specified directory.
|
static String[] |
ConfigurationUtils.parseLocalStateDirectories(Configuration configuration)
Extracts the local state directories as defined by
CheckpointingOptions.LOCAL_RECOVERY_TASK_MANAGER_STATE_ROOT_DIRS . |
static String[] |
ConfigurationUtils.parseTempDirectories(Configuration configuration)
Extracts the task manager directories for temporary files as defined by
CoreOptions.TMP_DIRS . |
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 |
---|---|
void |
HiveTableInputFormat.configure(Configuration parameters) |
void |
HiveTableOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
protected DispatcherResourceManagerComponentFactory<?> |
StandaloneJobClusterEntryPoint.createDispatcherResourceManagerComponentFactory(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
DefaultConfigurableOptionsFactory |
DefaultConfigurableOptionsFactory.configure(Configuration configuration)
Creates a
DefaultConfigurableOptionsFactory instance from a Configuration . |
OptionsFactory |
ConfigurableOptionsFactory.configure(Configuration configuration)
Creates a variant of the options factory that applies additional configuration parameters.
|
RocksDBStateBackend |
RocksDBStateBackend.configure(Configuration config,
ClassLoader classLoader)
Creates a copy of this state backend that uses the values defined in the configuration
for fields where that were not yet specified in this state backend.
|
RocksDBStateBackend |
RocksDBStateBackendFactory.createFromConfig(Configuration config,
ClassLoader classLoader) |
static RocksDBNativeMetricOptions |
RocksDBNativeMetricOptions.fromConfig(Configuration config)
Creates a
RocksDBNativeMetricOptions based on an
external configuration. |
Modifier and Type | Method and Description |
---|---|
void |
PluginFileSystemFactory.configure(Configuration config) |
void |
ConnectionLimitingFactory.configure(Configuration config) |
static FileSystemFactory |
ConnectionLimitingFactory.decorateIfLimited(FileSystemFactory factory,
String scheme,
Configuration config)
Decorates the given factory for a
ConnectionLimitingFactory , if the given
configuration configured connection limiting for the given file system scheme. |
static LimitedConnectionsFileSystem.ConnectionLimitingSettings |
LimitedConnectionsFileSystem.ConnectionLimitingSettings.fromConfig(Configuration config,
String fsScheme)
Parses and returns the settings for connection limiting, for the file system with
the given file system scheme.
|
static void |
FileSystem.initialize(Configuration config)
Deprecated.
use
FileSystem.initialize(Configuration, PluginManager) instead. |
static void |
FileSystem.initialize(Configuration config,
PluginManager pluginManager)
Initializes the shared file system settings.
|
Modifier and Type | Method and Description |
---|---|
default void |
Plugin.configure(Configuration config)
Optional method for plugins to pick up settings from the configuration.
|
static PluginManager |
PluginUtils.createPluginManagerFromRootFolder(Configuration configuration) |
static PluginConfig |
PluginConfig.fromConfiguration(Configuration configuration) |
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 |
ParquetInputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractAzureFSFactory.configure(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
SwiftFileSystemFactory.configure(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
OSSFileSystemFactory.configure(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractS3FileSystemFactory.configure(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
AnalyticHelper.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopReduceCombineFunction.open(Configuration parameters) |
void |
HadoopReduceFunction.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 |
---|---|
protected DispatcherResourceManagerComponentFactory<?> |
MesosJobClusterEntrypoint.createDispatcherResourceManagerComponentFactory(Configuration configuration) |
protected DispatcherResourceManagerComponentFactory<?> |
MesosSessionClusterEntrypoint.createDispatcherResourceManagerComponentFactory(Configuration configuration) |
protected void |
MesosJobClusterEntrypoint.initializeServices(Configuration config) |
protected void |
MesosSessionClusterEntrypoint.initializeServices(Configuration config) |
Constructor and Description |
---|
MesosJobClusterEntrypoint(Configuration config) |
MesosSessionClusterEntrypoint(Configuration config) |
Modifier and Type | Method and Description |
---|---|
static MesosTaskManagerParameters |
MesosTaskManagerParameters.create(Configuration flinkConfig)
Create the Mesos TaskManager parameters.
|
ResourceManager<RegisteredMesosWorkerNode> |
MesosResourceManagerFactory.createActiveResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Constructor and Description |
---|
MesosResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
Configuration flinkConfig,
MesosServices mesosServices,
MesosConfiguration mesosConfig,
MesosTaskManagerParameters taskManagerParameters,
ContainerSpecification taskManagerContainerSpec,
String webUiUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
static MesosServices |
MesosServicesUtils.createMesosServices(Configuration configuration,
String hostname)
Creates a
MesosServices instance depending on the high availability settings. |
MesosWorkerStore |
MesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor)
Creates a
MesosWorkerStore which is used to persist mesos worker in high availability
mode. |
MesosWorkerStore |
StandaloneMesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor) |
MesosWorkerStore |
ZooKeeperMesosServices.createMesosWorkerStore(Configuration configuration,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
MesosUtils.loadConfiguration(Configuration dynamicProperties,
org.slf4j.Logger log)
Loads the global configuration, adds the given dynamic properties configuration, and sets
the temp directory paths.
|
Modifier and Type | Method and Description |
---|---|
static void |
MesosUtils.applyOverlays(Configuration configuration,
ContainerSpecification containerSpec)
Generate a container specification as a TaskManager template.
|
static ContainerSpecification |
MesosUtils.createContainerSpec(Configuration flinkConfiguration) |
static MesosConfiguration |
MesosUtils.createMesosSchedulerConfiguration(Configuration flinkConfig,
String hostname)
Loads and validates the Mesos scheduler configuration.
|
static MesosTaskManagerParameters |
MesosUtils.createTmParameters(Configuration configuration,
org.slf4j.Logger log) |
static Configuration |
MesosUtils.loadConfiguration(Configuration dynamicProperties,
org.slf4j.Logger log)
Loads the global configuration, adds the given dynamic properties configuration, and sets
the temp directory paths.
|
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 | Field and Description |
---|---|
protected Configuration |
AbstractBlobCache.blobClientConfig
Configuration for the blob client like ssl parameters required to connect to the blob
server.
|
Modifier and Type | Method and Description |
---|---|
static BlobStoreService |
BlobUtils.createBlobStoreFromConfig(Configuration config)
Creates a BlobStore based on the parameters set in the configuration.
|
static List<PermanentBlobKey> |
BlobClient.uploadFiles(InetSocketAddress serverAddress,
Configuration clientConfig,
JobID jobId,
List<Path> files)
Uploads the JAR files to the
PermanentBlobService of the BlobServer at the
given address with HA as configured. |
Constructor and Description |
---|
AbstractBlobCache(Configuration blobClientConfig,
BlobView blobView,
org.slf4j.Logger logger,
InetSocketAddress serverAddress) |
BlobCacheService(Configuration blobClientConfig,
BlobView blobView,
InetSocketAddress serverAddress)
Instantiates a new BLOB cache.
|
BlobClient(InetSocketAddress serverAddress,
Configuration clientConfig)
Instantiates a new BLOB client.
|
BlobServer(Configuration config,
BlobStore blobStore)
Instantiates a new BLOB server and binds it to a free network port.
|
PermanentBlobCache(Configuration blobClientConfig,
BlobView blobView,
InetSocketAddress serverAddress)
Instantiates a new cache for permanent BLOBs which are also available in an HA store.
|
TransientBlobCache(Configuration blobClientConfig,
InetSocketAddress serverAddress)
Instantiates a new BLOB cache.
|
Modifier and Type | Method and Description |
---|---|
static void |
Checkpoints.disposeSavepoint(String pointer,
Configuration configuration,
ClassLoader classLoader,
org.slf4j.Logger logger) |
static StateBackend |
Checkpoints.loadStateBackend(Configuration configuration,
ClassLoader classLoader,
org.slf4j.Logger logger) |
Constructor and Description |
---|
ZooKeeperCheckpointRecoveryFactory(org.apache.curator.framework.CuratorFramework client,
Configuration config,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
BootstrapTools.cloneConfiguration(Configuration configuration)
Clones the given configuration and resets instance specific config options.
|
static Configuration |
BootstrapTools.generateTaskManagerConfiguration(Configuration baseConfig,
String jobManagerHostname,
int jobManagerPort,
int numSlots,
scala.concurrent.duration.FiniteDuration registrationTimeout)
Generate a task manager configuration.
|
Configuration |
ContainerSpecification.getFlinkConfiguration()
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 long |
ContaineredTaskManagerParameters.calculateCutoffMB(Configuration config,
long containerMemoryMB)
Calcuate cutoff memory size used by container, it will throw an
IllegalArgumentException
if the config is invalid or return the cutoff value if valid. |
static Configuration |
BootstrapTools.cloneConfiguration(Configuration configuration)
Clones the given configuration and resets instance specific config options.
|
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 ContainerSpecification |
ContainerSpecification.from(Configuration flinkConfiguration) |
static BootstrapTools.ForkJoinExecutorConfiguration |
BootstrapTools.ForkJoinExecutorConfiguration.fromConfiguration(Configuration configuration) |
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,
boolean hasKrb5,
Class<?> mainClass,
String mainArgs)
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,
int listeningPort,
org.slf4j.Logger logger,
BootstrapTools.ActorSystemExecutorConfiguration actorSystemExecutorConfiguration)
Starts an Actor System at a specific port.
|
static akka.actor.ActorSystem |
BootstrapTools.startActorSystem(Configuration configuration,
String actorSystemName,
String listeningAddress,
int listeningPort,
org.slf4j.Logger logger,
BootstrapTools.ActorSystemExecutorConfiguration actorSystemExecutorConfiguration)
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.ActorSystem |
BootstrapTools.startActorSystem(Configuration configuration,
String listeningAddress,
String portRangeDefinition,
org.slf4j.Logger logger,
BootstrapTools.ActorSystemExecutorConfiguration actorSystemExecutorConfiguration)
Starts an ActorSystem with the given configuration listening at the address/ports.
|
static akka.actor.ActorSystem |
BootstrapTools.startActorSystem(Configuration configuration,
String actorSystemName,
String listeningAddress,
String portRangeDefinition,
org.slf4j.Logger logger,
BootstrapTools.ActorSystemExecutorConfiguration actorSystemExecutorConfiguration)
Starts an ActorSystem with the given configuration listening at the address/ports.
|
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 a new config key to the value of a deprecated config key.
|
static void |
BootstrapTools.updateTmpDirectoriesInConfiguration(Configuration configuration,
String defaultDirs)
Set temporary configuration directories if necessary.
|
static void |
BootstrapTools.writeConfiguration(Configuration cfg,
File file)
Writes a Flink YAML config file from a Flink Configuration object.
|
Modifier and Type | Method and Description |
---|---|
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.
|
SSLStoreOverlay.Builder |
SSLStoreOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment (and global configuration).
|
KeytabOverlay.Builder |
KeytabOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment (and global configuration).
|
FlinkDistributionOverlay.Builder |
FlinkDistributionOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment.
|
Krb5ConfOverlay.Builder |
Krb5ConfOverlay.Builder.fromEnvironment(Configuration globalConfiguration)
Configures the overlay using the current environment.
|
Modifier and Type | Method and Description |
---|---|
Dispatcher |
SessionDispatcherFactory.createDispatcher(Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricQueryServiceAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist) |
T |
DispatcherFactory.createDispatcher(Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricQueryServiceAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist)
Create a
Dispatcher of the given type T . |
MiniDispatcher |
JobDispatcherFactory.createDispatcher(Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricQueryServiceAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist) |
static HistoryServerArchivist |
HistoryServerArchivist.createHistoryServerArchivist(Configuration configuration,
JsonArchivist jsonArchivist,
Executor ioExecutor) |
JobManagerRunner |
DefaultJobManagerRunnerFactory.createJobManagerRunner(JobGraph jobGraph,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
JobManagerSharedServices jobManagerServices,
JobManagerJobMetricGroupFactory jobManagerJobMetricGroupFactory,
FatalErrorHandler fatalErrorHandler) |
JobManagerRunner |
JobManagerRunnerFactory.createJobManagerRunner(JobGraph jobGraph,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
JobManagerSharedServices jobManagerServices,
JobManagerJobMetricGroupFactory jobManagerJobMetricGroupFactory,
FatalErrorHandler fatalErrorHandler) |
Constructor and Description |
---|
Dispatcher(RpcService rpcService,
String endpointId,
Configuration configuration,
HighAvailabilityServices highAvailabilityServices,
SubmittedJobGraphStore submittedJobGraphStore,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricServiceQueryAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
JobManagerRunnerFactory jobManagerRunnerFactory,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist) |
DispatcherRestEndpoint(RestServerEndpointConfiguration endpointConfiguration,
GatewayRetriever<DispatcherGateway> leaderRetriever,
Configuration clusterConfiguration,
RestHandlerConfiguration restConfiguration,
GatewayRetriever<ResourceManagerGateway> resourceManagerRetriever,
TransientBlobService transientBlobService,
ScheduledExecutorService executor,
MetricFetcher metricFetcher,
LeaderElectionService leaderElectionService,
ExecutionGraphCache executionGraphCache,
FatalErrorHandler fatalErrorHandler) |
MiniDispatcher(RpcService rpcService,
String endpointId,
Configuration configuration,
HighAvailabilityServices highAvailabilityServices,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricQueryServiceAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
JobManagerRunnerFactory jobManagerRunnerFactory,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist,
JobGraph jobGraph,
ClusterEntrypoint.ExecutionMode executionMode) |
StandaloneDispatcher(RpcService rpcService,
String endpointId,
Configuration configuration,
HighAvailabilityServices highAvailabilityServices,
GatewayRetriever<ResourceManagerGateway> resourceManagerGatewayRetriever,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
JobManagerMetricGroup jobManagerMetricGroup,
String metricQueryServiceAddress,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
JobManagerRunnerFactory jobManagerRunnerFactory,
FatalErrorHandler fatalErrorHandler,
HistoryServerArchivist historyServerArchivist) |
Modifier and Type | Method and Description |
---|---|
protected static Configuration |
ClusterEntrypoint.loadConfiguration(EntrypointClusterConfiguration entrypointClusterConfiguration) |
Constructor and Description |
---|
ClusterEntrypoint(Configuration configuration) |
JobClusterEntrypoint(Configuration configuration) |
SessionClusterEntrypoint(Configuration configuration) |
StandaloneSessionClusterEntrypoint(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
DispatcherResourceManagerComponent<T> |
DispatcherResourceManagerComponentFactory.create(Configuration configuration,
Executor ioExecutor,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
MetricQueryServiceRetriever metricQueryServiceRetriever,
FatalErrorHandler fatalErrorHandler) |
DispatcherResourceManagerComponent<T> |
AbstractDispatcherResourceManagerComponentFactory.create(Configuration configuration,
Executor ioExecutor,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
ArchivedExecutionGraphStore archivedExecutionGraphStore,
MetricQueryServiceRetriever metricQueryServiceRetriever,
FatalErrorHandler fatalErrorHandler) |
static FileJobGraphRetriever |
FileJobGraphRetriever.createFrom(Configuration configuration) |
JobGraph |
FileJobGraphRetriever.retrieveJobGraph(Configuration configuration) |
JobGraph |
JobGraphRetriever.retrieveJobGraph(Configuration configuration)
Retrieve the
JobGraph . |
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,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time rpcTimeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
BlobWriter blobWriter,
Time allocationTimeout,
org.slf4j.Logger log,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker)
Builds the ExecutionGraph from the JobGraph.
|
static ExecutionGraph |
ExecutionGraphBuilder.buildGraph(ExecutionGraph prior,
JobGraph jobGraph,
Configuration jobManagerConfig,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time rpcTimeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
BlobWriter blobWriter,
Time allocationTimeout,
org.slf4j.Logger log,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker,
FailoverStrategy.Factory failoverStrategyFactory) |
Constructor and Description |
---|
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<PermanentBlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
TaskInformation(JobVertexID jobVertexId,
String taskName,
int numberOfSubtasks,
int maxNumberOfSubtaks,
String invokableClassName,
Configuration taskConfiguration) |
Modifier and Type | Method and Description |
---|---|
static FailoverStrategy.Factory |
FailoverStrategyLoader.loadFailoverStrategy(Configuration config,
org.slf4j.Logger logger)
Loads a FailoverStrategy Factory from the given configuration.
|
Modifier and Type | Method and Description |
---|---|
static FailureRateRestartBackoffTimeStrategy.FailureRateRestartBackoffTimeStrategyFactory |
FailureRateRestartBackoffTimeStrategy.createFactory(Configuration configuration) |
static FixedDelayRestartBackoffTimeStrategy.FixedDelayRestartBackoffTimeStrategyFactory |
FixedDelayRestartBackoffTimeStrategy.createFactory(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static PartitionReleaseStrategy.Factory |
PartitionReleaseStrategyFactoryLoader.loadPartitionReleaseStrategyFactory(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static FailureRateRestartStrategy.FailureRateRestartStrategyFactory |
FailureRateRestartStrategy.createFactory(Configuration configuration) |
static NoRestartStrategy.NoRestartStrategyFactory |
NoRestartStrategy.createFactory(Configuration configuration)
Creates a NoRestartStrategyFactory instance.
|
static FixedDelayRestartStrategy.FixedDelayRestartStrategyFactory |
FixedDelayRestartStrategy.createFactory(Configuration configuration)
Creates a FixedDelayRestartStrategy from the given Configuration.
|
static RestartStrategyFactory |
RestartStrategyFactory.createRestartStrategyFactory(Configuration configuration)
Creates a
RestartStrategy instance from the given Configuration . |
Modifier and Type | Method and Description |
---|---|
void |
HadoopFsFactory.configure(Configuration config) |
Modifier and Type | Method and Description |
---|---|
static HeartbeatServices |
HeartbeatServices.fromConfiguration(Configuration configuration)
Creates an HeartbeatServices instance from a
Configuration . |
Modifier and Type | Method and Description |
---|---|
static HighAvailabilityServices |
HighAvailabilityServicesUtils.createAvailableOrEmbeddedServices(Configuration config,
Executor executor) |
HighAvailabilityServices |
HighAvailabilityServicesFactory.createHAServices(Configuration configuration,
Executor executor)
Creates an
HighAvailabilityServices instance. |
static HighAvailabilityServices |
HighAvailabilityServicesUtils.createHighAvailabilityServices(Configuration configuration,
Executor executor,
HighAvailabilityServicesUtils.AddressResolution addressResolution) |
static Tuple2<String,Integer> |
HighAvailabilityServicesUtils.getJobManagerAddress(Configuration configuration)
Returns the JobManager's hostname and port extracted from the given
Configuration . |
Constructor and Description |
---|
ZooKeeperHaServices(org.apache.curator.framework.CuratorFramework client,
Executor executor,
Configuration configuration,
BlobStoreService blobStoreService) |
ZooKeeperRunningJobsRegistry(org.apache.curator.framework.CuratorFramework client,
Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
NettyShuffleMaster |
NettyShuffleServiceFactory.createShuffleMaster(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
NettyConfig.getConfig() |
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.
|
Configuration |
InputOutputFormatContainer.getParameters(OperatorID operatorId) |
Modifier and Type | Method and Description |
---|---|
InputOutputFormatContainer |
InputOutputFormatContainer.addParameters(OperatorID operatorId,
Configuration parameters) |
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 |
---|---|
static HighAvailabilityMode |
HighAvailabilityMode.fromConfig(Configuration config)
Return the configured
HighAvailabilityMode . |
static boolean |
HighAvailabilityMode.isHighAvailabilityModeActivated(Configuration configuration)
Returns true if the defined recovery mode supports high availability.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
JobMasterConfiguration.getConfiguration() |
Modifier and Type | Method and Description |
---|---|
static JobMasterConfiguration |
JobMasterConfiguration.fromConfiguration(Configuration configuration) |
static JobManagerSharedServices |
JobManagerSharedServices.fromConfiguration(Configuration config,
BlobServer blobServer) |
Constructor and Description |
---|
JobMasterConfiguration(Time rpcTimeout,
Time slotRequestTimeout,
String tmpDirectory,
RetryingRegistrationConfiguration retryingRegistrationConfiguration,
Configuration configuration) |
MiniDispatcherRestEndpoint(RestServerEndpointConfiguration endpointConfiguration,
GatewayRetriever<? extends RestfulGateway> leaderRetriever,
Configuration clusterConfiguration,
RestHandlerConfiguration restConfiguration,
GatewayRetriever<ResourceManagerGateway> resourceManagerRetriever,
TransientBlobService transientBlobService,
ScheduledExecutorService executor,
MetricFetcher metricFetcher,
LeaderElectionService leaderElectionService,
ExecutionGraphCache executionGraphCache,
FatalErrorHandler fatalErrorHandler) |
Modifier and Type | Method and Description |
---|---|
static DefaultSchedulerFactory |
DefaultSchedulerFactory.fromConfiguration(Configuration configuration) |
static DefaultSlotPoolFactory |
DefaultSlotPoolFactory.fromConfiguration(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static List<ReporterSetup> |
ReporterSetup.fromConfiguration(Configuration configuration) |
static MetricRegistryConfiguration |
MetricRegistryConfiguration.fromConfiguration(Configuration configuration)
Create a metric registry configuration object from the given
Configuration . |
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 |
---|---|
static RpcService |
MetricUtils.startMetricsRpcService(Configuration configuration,
String hostname) |
Modifier and Type | Method and Description |
---|---|
protected Collection<? extends DispatcherResourceManagerComponent<?>> |
MiniCluster.createDispatcherResourceManagerComponents(Configuration configuration,
MiniCluster.RpcServiceFactory rpcServiceFactory,
HighAvailabilityServices haServices,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
MetricQueryServiceRetriever metricQueryServiceRetriever,
FatalErrorHandler fatalErrorHandler) |
protected HighAvailabilityServices |
MiniCluster.createHighAvailabilityServices(Configuration configuration,
Executor executor) |
protected MetricRegistryImpl |
MiniCluster.createMetricRegistry(Configuration config)
Factory method to create the metric registry for the mini cluster.
|
MiniClusterConfiguration.Builder |
MiniClusterConfiguration.Builder.setConfiguration(Configuration configuration1) |
Constructor and Description |
---|
MiniClusterConfiguration(Configuration configuration,
int numTaskManagers,
RpcServiceSharing rpcServiceSharing,
String commonBindAddress) |
Modifier and Type | Method and Description |
---|---|
static SSLHandlerFactory |
SSLUtils.createInternalClientSSLEngineFactory(Configuration config)
Creates a SSLEngineFactory to be used by internal communication client endpoints.
|
static SSLHandlerFactory |
SSLUtils.createInternalServerSSLEngineFactory(Configuration config)
Creates a SSLEngineFactory to be used by internal communication server endpoints.
|
static SSLHandlerFactory |
SSLUtils.createRestClientSSLEngineFactory(Configuration config)
Creates a
SSLHandlerFactory to be used by the REST Clients. |
static org.apache.flink.shaded.netty4.io.netty.handler.ssl.SslContext |
SSLUtils.createRestNettySSLContext(Configuration config,
boolean clientMode,
org.apache.flink.shaded.netty4.io.netty.handler.ssl.ClientAuth clientAuth,
org.apache.flink.shaded.netty4.io.netty.handler.ssl.SslProvider provider)
Creates an SSL context for the external REST SSL.
|
static SSLHandlerFactory |
SSLUtils.createRestServerSSLEngineFactory(Configuration config)
Creates a
SSLHandlerFactory to be used by the REST Servers. |
static SSLContext |
SSLUtils.createRestSSLContext(Configuration config,
boolean clientMode)
Creates an SSL context for clients against the external REST endpoint.
|
static SocketFactory |
SSLUtils.createSSLClientSocketFactory(Configuration config)
Creates a factory for SSL Client Sockets from the given configuration.
|
static ServerSocketFactory |
SSLUtils.createSSLServerSocketFactory(Configuration config)
Creates a factory for SSL Server Sockets from the given configuration.
|
static boolean |
SSLUtils.isInternalSSLEnabled(Configuration sslConfig)
Checks whether SSL for internal communication (rpc, data transport, blob server) is enabled.
|
static boolean |
SSLUtils.isRestSSLAuthenticationEnabled(Configuration sslConfig)
Checks whether mutual SSL authentication for the external REST endpoint is enabled.
|
static boolean |
SSLUtils.isRestSSLEnabled(Configuration sslConfig)
Checks whether SSL for the external REST endpoint is enabled.
|
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 |
---|---|
static RetryingRegistrationConfiguration |
RetryingRegistrationConfiguration.fromConfiguration(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
ActiveResourceManagerFactory.createActiveResourceManagerConfiguration(Configuration originalConfiguration) |
Modifier and Type | Method and Description |
---|---|
protected abstract ResourceManager<T> |
ActiveResourceManagerFactory.createActiveResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
static Configuration |
ActiveResourceManagerFactory.createActiveResourceManagerConfiguration(Configuration originalConfiguration) |
ResourceManager<ResourceID> |
StandaloneResourceManagerFactory.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
ResourceManager<T> |
ResourceManagerFactory.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
ResourceManager<T> |
ActiveResourceManagerFactory.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
static Collection<ResourceProfile> |
ResourceManager.createWorkerSlotProfiles(Configuration config) |
static ResourceManagerRuntimeServicesConfiguration |
ResourceManagerRuntimeServicesConfiguration.fromConfiguration(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static SlotManagerConfiguration |
SlotManagerConfiguration.fromConfiguration(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static RestHandlerConfiguration |
RestHandlerConfiguration.fromConfiguration(Configuration configuration) |
Constructor and Description |
---|
ClusterConfigHandler(GatewayRetriever<? extends RestfulGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
MessageHeaders<EmptyRequestBody,ClusterConfigurationInfo,EmptyMessageParameters> messageHeaders,
Configuration configuration) |
Constructor and Description |
---|
JobSubmitHandler(GatewayRetriever<? extends DispatcherGateway> leaderRetriever,
Time timeout,
Map<String,String> headers,
Executor executor,
Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static <T extends RestfulGateway> |
MetricFetcherImpl.fromConfiguration(Configuration configuration,
MetricQueryServiceRetriever metricQueryServiceGatewayRetriever,
GatewayRetriever<T> dispatcherGatewayRetriever,
ExecutorService executor) |
Modifier and Type | Method and Description |
---|---|
static ClusterConfigurationInfo |
ClusterConfigurationInfo.from(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
AkkaRpcServiceConfiguration.getConfiguration() |
Modifier and Type | Method and Description |
---|---|
static RpcService |
AkkaRpcServiceUtils.createRpcService(String hostname,
int port,
Configuration configuration)
Utility method to create RPC service from configuration and hostname, port.
|
static RpcService |
AkkaRpcServiceUtils.createRpcService(String hostname,
String portRangeDefinition,
Configuration configuration)
Utility method to create RPC service from configuration and hostname, port.
|
static RpcService |
AkkaRpcServiceUtils.createRpcService(String hostname,
String portRangeDefinition,
Configuration configuration,
String actorSystemName,
BootstrapTools.ActorSystemExecutorConfiguration actorSystemExecutorConfiguration)
Utility method to create RPC service from configuration and hostname, port.
|
static long |
AkkaRpcServiceUtils.extractMaximumFramesize(Configuration configuration) |
static AkkaRpcServiceConfiguration |
AkkaRpcServiceConfiguration.fromConfiguration(Configuration configuration) |
static String |
AkkaRpcServiceUtils.getRpcUrl(String hostname,
int port,
String endpointName,
HighAvailabilityServicesUtils.AddressResolution addressResolution,
Configuration config) |
Constructor and Description |
---|
AkkaRpcServiceConfiguration(Configuration configuration,
Time timeout,
long maximumFramesize) |
Modifier and Type | Method and Description |
---|---|
SchedulerNG |
DefaultSchedulerFactory.createInstance(org.slf4j.Logger log,
JobGraph jobGraph,
BackPressureStatsTracker backPressureStatsTracker,
Executor ioExecutor,
Configuration jobMasterConfiguration,
SlotProvider slotProvider,
ScheduledExecutorService futureExecutor,
ClassLoader userCodeLoader,
CheckpointRecoveryFactory checkpointRecoveryFactory,
Time rpcTimeout,
BlobWriter blobWriter,
JobManagerJobMetricGroup jobManagerJobMetricGroup,
Time slotRequestTimeout,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker) |
SchedulerNG |
SchedulerNGFactory.createInstance(org.slf4j.Logger log,
JobGraph jobGraph,
BackPressureStatsTracker backPressureStatsTracker,
Executor ioExecutor,
Configuration jobMasterConfiguration,
SlotProvider slotProvider,
ScheduledExecutorService futureExecutor,
ClassLoader userCodeLoader,
CheckpointRecoveryFactory checkpointRecoveryFactory,
Time rpcTimeout,
BlobWriter blobWriter,
JobManagerJobMetricGroup jobManagerJobMetricGroup,
Time slotRequestTimeout,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker) |
SchedulerNG |
LegacySchedulerFactory.createInstance(org.slf4j.Logger log,
JobGraph jobGraph,
BackPressureStatsTracker backPressureStatsTracker,
Executor ioExecutor,
Configuration jobMasterConfiguration,
SlotProvider slotProvider,
ScheduledExecutorService futureExecutor,
ClassLoader userCodeLoader,
CheckpointRecoveryFactory checkpointRecoveryFactory,
Time rpcTimeout,
BlobWriter blobWriter,
JobManagerJobMetricGroup jobManagerJobMetricGroup,
Time slotRequestTimeout,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker) |
Constructor and Description |
---|
DefaultScheduler(org.slf4j.Logger log,
JobGraph jobGraph,
BackPressureStatsTracker backPressureStatsTracker,
Executor ioExecutor,
Configuration jobMasterConfiguration,
SlotProvider slotProvider,
ScheduledExecutorService futureExecutor,
ClassLoader userCodeLoader,
CheckpointRecoveryFactory checkpointRecoveryFactory,
Time rpcTimeout,
BlobWriter blobWriter,
JobManagerJobMetricGroup jobManagerJobMetricGroup,
Time slotRequestTimeout,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker) |
LegacyScheduler(org.slf4j.Logger log,
JobGraph jobGraph,
BackPressureStatsTracker backPressureStatsTracker,
Executor ioExecutor,
Configuration jobMasterConfiguration,
SlotProvider slotProvider,
ScheduledExecutorService futureExecutor,
ClassLoader userCodeLoader,
CheckpointRecoveryFactory checkpointRecoveryFactory,
Time rpcTimeout,
RestartStrategyFactory restartStrategyFactory,
BlobWriter blobWriter,
JobManagerJobMetricGroup jobManagerJobMetricGroup,
Time slotRequestTimeout,
ShuffleMaster<?> shuffleMaster,
PartitionTracker partitionTracker) |
Modifier and Type | Method and Description |
---|---|
Configuration |
SecurityConfiguration.getFlinkConfig() |
Constructor and Description |
---|
SecurityConfiguration(Configuration flinkConf)
Create a security configuration from the global configuration.
|
SecurityConfiguration(Configuration flinkConf,
List<SecurityModuleFactory> securityModuleFactories)
Create a security configuration from the global configuration.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
ShuffleEnvironmentContext.getConfiguration() |
Modifier and Type | Method and Description |
---|---|
ShuffleMaster<SD> |
ShuffleServiceFactory.createShuffleMaster(Configuration configuration)
Factory method to create a specific
ShuffleMaster implementation. |
static ShuffleServiceFactory<?,?,?> |
ShuffleServiceLoader.loadShuffleServiceFactory(Configuration configuration) |
Constructor and Description |
---|
ShuffleEnvironmentContext(Configuration configuration,
ResourceID taskExecutorResourceId,
long maxJvmHeapMemory,
boolean localCommunicationOnly,
InetAddress hostAddress,
TaskEventPublisher eventPublisher,
MetricGroup parentMetricGroup) |
Modifier and Type | Method and Description |
---|---|
StateBackend |
ConfigurableStateBackend.configure(Configuration config,
ClassLoader classLoader)
Creates a variant of the state backend that applies additional configuration parameters.
|
T |
StateBackendFactory.createFromConfig(Configuration config,
ClassLoader classLoader)
Creates the state backend, optionally using the given configuration.
|
static StateBackend |
StateBackendLoader.fromApplicationOrConfigOrDefault(StateBackend fromApplication,
Configuration config,
ClassLoader classLoader,
org.slf4j.Logger logger)
Checks if an application-defined state backend is given, and if not, loads the state
backend from the configuration, from the parameter 'state.backend', as defined
in
CheckpointingOptions.STATE_BACKEND . |
static StateBackend |
StateBackendLoader.loadStateBackendFromConfig(Configuration config,
ClassLoader classLoader,
org.slf4j.Logger logger)
Loads the state backend from the configuration, from the parameter 'state.backend', as defined
in
CheckpointingOptions.STATE_BACKEND . |
Modifier and Type | Method and Description |
---|---|
FsStateBackend |
FsStateBackend.configure(Configuration config,
ClassLoader classLoader)
Creates a copy of this state backend that uses the values defined in the configuration
for fields where that were not specified in this state backend.
|
FsStateBackend |
FsStateBackendFactory.createFromConfig(Configuration config,
ClassLoader classLoader) |
Constructor and Description |
---|
AbstractFileStateBackend(Path baseCheckpointPath,
Path baseSavepointPath,
Configuration configuration)
Creates a new backend using the given checkpoint-/savepoint directories, or the values defined in
the given configuration.
|
Modifier and Type | Method and Description |
---|---|
MemoryStateBackend |
MemoryStateBackend.configure(Configuration config,
ClassLoader classLoader)
Creates a copy of this state backend that uses the values defined in the configuration
for fields where that were not specified in this state backend.
|
MemoryStateBackend |
MemoryStateBackendFactory.createFromConfig(Configuration config,
ClassLoader classLoader) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskManagerConfiguration.getConfiguration() |
Configuration |
TaskManagerServicesConfiguration.getConfiguration() |
Modifier and Type | Method and Description |
---|---|
static long |
TaskManagerServices.calculateHeapSizeMB(long totalJavaMemorySizeMB,
Configuration config)
Calculates the amount of heap memory to use (to set via -Xmx and -Xms)
based on the total memory to use and the given configuration parameters.
|
static RpcService |
TaskManagerRunner.createRpcService(Configuration configuration,
HighAvailabilityServices haServices)
Create a RPC service for the task manager.
|
static TaskManagerConfiguration |
TaskManagerConfiguration.fromConfiguration(Configuration configuration) |
static QueryableStateConfiguration |
QueryableStateConfiguration.fromConfiguration(Configuration config)
Creates the
QueryableStateConfiguration from the given Configuration. |
static TaskManagerServicesConfiguration |
TaskManagerServicesConfiguration.fromConfiguration(Configuration configuration,
ResourceID resourceID,
InetAddress remoteAddress,
long freeHeapMemoryWithDefrag,
long maxJvmHeapMemory,
boolean localCommunicationOnly)
Utility method to extract TaskManager config parameters from the configuration and to
sanity check them.
|
static long |
TaskManagerServices.getManagedMemoryFromHeapAndManaged(Configuration config,
long heapAndManagedMemory)
Gets the size of managed memory from the heap size after subtracting network buffer memory.
|
static long |
TaskManagerServices.getManagedMemoryFromProcessMemory(Configuration config,
long totalProcessMemory)
Gets the size of managed memory from the JVM process size, which at that point includes
network buffer memory, managed memory, and non-flink-managed heap memory.
|
static long |
TaskManagerServices.getReservedNetworkMemory(Configuration config,
long totalProcessMemory)
Gets the amount of memory reserved for networking, given the total JVM memory.
|
static void |
TaskManagerRunner.runTaskManager(Configuration configuration,
ResourceID resourceId) |
static TaskExecutor |
TaskManagerRunner.startTaskManager(Configuration configuration,
ResourceID resourceID,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
BlobCacheService blobCacheService,
boolean localCommunicationOnly,
FatalErrorHandler fatalErrorHandler) |
Constructor and Description |
---|
TaskManagerConfiguration(int numberSlots,
String[] tmpDirectories,
Time timeout,
Time maxRegistrationDuration,
Time initialRegistrationPause,
Time maxRegistrationPause,
Time refusedRegistrationPause,
Configuration configuration,
boolean exitJvmOnOutOfMemory,
FlinkUserCodeClassLoaders.ResolveOrder classLoaderResolveOrder,
String[] alwaysParentFirstLoaderPatterns,
String taskManagerLogPath,
String taskManagerStdoutPath,
RetryingRegistrationConfiguration retryingRegistrationConfiguration) |
TaskManagerRunner(Configuration configuration,
ResourceID resourceId) |
TaskManagerServicesConfiguration(Configuration configuration,
ResourceID resourceID,
InetAddress taskManagerAddress,
boolean localCommunicationOnly,
String[] tmpDirPaths,
String[] localRecoveryStateRootDirectories,
long freeHeapMemoryWithDefrag,
long maxJvmHeapMemory,
boolean localRecoveryEnabled,
QueryableStateConfiguration queryableStateConfig,
int numberOfSlots,
long configuredMemory,
MemoryType memoryType,
boolean preAllocateMemory,
float memoryFraction,
int pageSize,
long timerServiceShutdownTimeout,
RetryingRegistrationConfiguration retryingRegistrationConfiguration,
Optional<Time> systemResourceMetricsProbingInterval) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskManagerRuntimeInfo.getConfiguration()
Gets the configuration that the TaskManager was started with.
|
Configuration |
RuntimeEnvironment.getJobConfiguration() |
Configuration |
Task.getJobConfiguration() |
Configuration |
RuntimeEnvironment.getTaskConfiguration() |
Configuration |
Task.getTaskConfiguration() |
Modifier and Type | Method and Description |
---|---|
static long |
NettyShuffleEnvironmentConfiguration.calculateNetworkBufferMemory(long totalJavaMemorySize,
Configuration config)
Calculates the amount of memory used for network buffers based on the total memory to use and
the according configuration parameters.
|
static long |
NettyShuffleEnvironmentConfiguration.calculateNewNetworkBufferMemory(Configuration config,
long maxJvmHeapMemory)
Calculates the amount of memory used for network buffers inside the current JVM instance
based on the available heap or the max heap size and the according configuration parameters.
|
static NettyShuffleEnvironmentConfiguration |
NettyShuffleEnvironmentConfiguration.fromConfiguration(Configuration configuration,
long maxJvmHeapMemory,
boolean localTaskManagerCommunication,
InetAddress taskManagerAddress)
Utility method to extract network related parameters from the configuration and to
sanity check them.
|
static boolean |
NettyShuffleEnvironmentConfiguration.hasNewNetworkConfig(Configuration config)
Returns whether the new network buffer memory configuration is present in the configuration
object, i.e.
|
static void |
MemoryLogger.startIfConfigured(org.slf4j.Logger logger,
Configuration configuration,
CompletableFuture<Void> taskManagerTerminationFuture) |
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,
TaskStateManager taskStateManager,
GlobalAggregateManager aggregateManager,
AccumulatorRegistry accumulatorRegistry,
TaskKvStateRegistry kvStateRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
TaskEventDispatcher taskEventDispatcher,
CheckpointResponder checkpointResponder,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
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(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperLeaderElectionService instance. |
static ZooKeeperLeaderElectionService |
ZooKeeperUtils.createLeaderElectionService(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
String pathSuffix)
Creates a
ZooKeeperLeaderElectionService instance. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
StandaloneLeaderRetrievalService from the given configuration. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration,
boolean resolveInitialHostName,
String jobManagerName)
Creates a
StandaloneLeaderRetrievalService form the given configuration and the
JobManager name. |
static ZooKeeperLeaderRetrievalService |
ZooKeeperUtils.createLeaderRetrievalService(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperLeaderRetrievalService instance. |
static ZooKeeperLeaderRetrievalService |
ZooKeeperUtils.createLeaderRetrievalService(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
String pathSuffix)
Creates a
ZooKeeperLeaderRetrievalService instance. |
static ZooKeeperSubmittedJobGraphStore |
ZooKeeperUtils.createSubmittedJobGraphs(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperSubmittedJobGraphStore instance. |
static ZooKeeperUtils.ZkClientACLMode |
ZooKeeperUtils.ZkClientACLMode.fromConfig(Configuration config)
Return the configured
ZooKeeperUtils.ZkClientACLMode . |
static Configuration |
HadoopUtils.getHadoopConfiguration(Configuration flinkConfiguration) |
static float |
ConfigurationParserUtils.getManagedMemoryFraction(Configuration configuration)
Parses the configuration to get the fraction of managed memory and validates the value.
|
static long |
ConfigurationParserUtils.getManagedMemorySize(Configuration configuration)
Parses the configuration to get the managed memory size and validates the value.
|
static MemoryType |
ConfigurationParserUtils.getMemoryType(Configuration configuration)
Parses the configuration to get the type of memory.
|
static int |
ConfigurationParserUtils.getPageSize(Configuration configuration)
Parses the configuration to get the page size and validates the value.
|
static int |
ConfigurationParserUtils.getSlot(Configuration configuration)
Parses the configuration to get the number of slots and validates the value.
|
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. |
void |
HadoopConfigLoader.setFlinkConfig(Configuration config) |
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 | Field and Description |
---|---|
protected Configuration |
WebMonitorEndpoint.clusterConfiguration |
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 WebMonitorExtension |
WebMonitorUtils.loadWebSubmissionExtension(GatewayRetriever<? extends DispatcherGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
CompletableFuture<String> localAddressFuture,
Path uploadDir,
Executor executor,
Configuration configuration)
Loads the
WebMonitorExtension which enables web submission. |
Constructor and Description |
---|
WebMonitorEndpoint(RestServerEndpointConfiguration endpointConfiguration,
GatewayRetriever<? extends T> leaderRetriever,
Configuration clusterConfiguration,
RestHandlerConfiguration restConfiguration,
GatewayRetriever<ResourceManagerGateway> resourceManagerRetriever,
TransientBlobService transientBlobService,
ScheduledExecutorService executor,
MetricFetcher metricFetcher,
LeaderElectionService leaderElectionService,
ExecutionGraphCache executionGraphCache,
FatalErrorHandler fatalErrorHandler) |
WebSubmissionExtension(Configuration configuration,
GatewayRetriever<? extends DispatcherGateway> leaderRetriever,
Map<String,String> responseHeaders,
CompletableFuture<String> localAddressFuture,
Path jarDir,
Executor executor,
Time timeout) |
Constructor and Description |
---|
JarListHandler(GatewayRetriever<? extends RestfulGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
MessageHeaders<EmptyRequestBody,JarListInfo,EmptyMessageParameters> messageHeaders,
CompletableFuture<String> localAddressFuture,
File jarDir,
Configuration configuration,
Executor executor) |
JarPlanHandler(GatewayRetriever<? extends RestfulGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
MessageHeaders<JarPlanRequestBody,JobPlanInfo,org.apache.flink.runtime.webmonitor.handlers.JarPlanMessageParameters> messageHeaders,
Path jarDir,
Configuration configuration,
Executor executor) |
JarPlanHandler(GatewayRetriever<? extends RestfulGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
MessageHeaders<JarPlanRequestBody,JobPlanInfo,org.apache.flink.runtime.webmonitor.handlers.JarPlanMessageParameters> messageHeaders,
Path jarDir,
Configuration configuration,
Executor executor,
java.util.function.Function<JobGraph,JobPlanInfo> planGenerator) |
JarRunHandler(GatewayRetriever<? extends DispatcherGateway> leaderRetriever,
Time timeout,
Map<String,String> responseHeaders,
MessageHeaders<JarRunRequestBody,JarRunResponseBody,JarRunMessageParameters> messageHeaders,
Path jarDir,
Configuration configuration,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
JobGraph |
JarHandlerUtils.JarHandlerContext.toJobGraph(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
static Optional<URL> |
HistoryServerUtils.getHistoryServerURL(Configuration configuration) |
static boolean |
HistoryServerUtils.isSSLEnabled(Configuration config) |
Constructor and Description |
---|
HistoryServer(Configuration config) |
HistoryServer(Configuration config,
CountDownLatch numArchivedJobs) |
Constructor and Description |
---|
WebFrontendBootstrap(Router router,
org.slf4j.Logger log,
File directory,
SSLHandlerFactory serverSSLFactory,
String configuredAddress,
int configuredPort,
Configuration config) |
Constructor and Description |
---|
ZooKeeperUtilityFactory(Configuration configuration,
String path) |
Modifier and Type | Method and Description |
---|---|
abstract void |
KeyedStateReaderFunction.open(Configuration parameters)
Initialization method for the function.
|
Modifier and Type | Method and Description |
---|---|
void |
KeyedStateInputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
SavepointOutputFormat.configure(Configuration parameters) |
void |
OperatorSubtaskStateReducer.open(Configuration parameters) |
void |
BoundedOneInputStreamTaskRunner.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
Configuration |
SavepointEnvironment.getJobConfiguration() |
Configuration |
SavepointEnvironment.getTaskConfiguration() |
Modifier and Type | Method and Description |
---|---|
SavepointEnvironment.Builder |
SavepointEnvironment.Builder.setConfiguration(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
RemoteStreamEnvironment.getClientConfiguration() |
protected Configuration |
LocalStreamEnvironment.getConfiguration() |
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 . |
static JobExecutionResult |
RemoteStreamEnvironment.executeRemotely(StreamExecutionEnvironment streamExecutionEnvironment,
List<URL> jarFiles,
String host,
int port,
Configuration clientConfiguration,
List<URL> globalClasspaths,
String jobName,
SavepointRestoreSettings savepointRestoreSettings)
Executes the job remotely.
|
Constructor and Description |
---|
LocalStreamEnvironment(Configuration configuration)
Creates a new mini cluster 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.
|
RemoteStreamEnvironment(String host,
int port,
Configuration clientConfiguration,
String[] jarFiles,
URL[] globalClasspaths,
SavepointRestoreSettings savepointRestoreSettings)
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 |
PrintSinkFunction.open(Configuration parameters) |
void |
SocketClientSink.open(Configuration parameters)
Initialize the connection with the Socket in the server.
|
void |
OutputFormatSinkFunction.open(Configuration parameters)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
void |
StreamingFileSink.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
InputFormatSourceFunction.open(Configuration parameters) |
void |
MultipleIdsMessageAcknowledgingSourceBase.open(Configuration parameters) |
void |
FromSplittableIteratorFunction.open(Configuration parameters) |
void |
ContinuousFileMonitoringFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
FoldApplyWindowFunction.open(Configuration configuration)
Deprecated.
|
void |
FoldApplyProcessWindowFunction.open(Configuration configuration)
Deprecated.
|
void |
FoldApplyAllWindowFunction.open(Configuration configuration)
Deprecated.
|
void |
FoldApplyProcessAllWindowFunction.open(Configuration configuration)
Deprecated.
|
void |
ReduceApplyProcessWindowFunction.open(Configuration configuration) |
void |
ReduceApplyProcessAllWindowFunction.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 |
CassandraSinkBase.open(Configuration configuration) |
void |
CassandraPojoSink.open(Configuration configuration) |
void |
AbstractCassandraTupleSink.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
void |
ElasticsearchSinkBase.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
static org.apache.hadoop.fs.FileSystem |
BucketingSink.createHadoopFileSystem(org.apache.hadoop.fs.Path path,
Configuration extraUserConf)
Deprecated.
|
void |
BucketingSink.open(Configuration parameters)
Deprecated.
|
BucketingSink<T> |
BucketingSink.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 |
PubSubSource.open(Configuration configuration) |
void |
PubSubSink.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
void |
FlinkKafkaProducer.open(Configuration configuration)
Initializes the connection to Kafka.
|
void |
FlinkKafkaProducer011.open(Configuration configuration)
Initializes the connection to Kafka.
|
void |
FlinkKafkaProducerBase.open(Configuration configuration)
Initializes the connection to Kafka.
|
void |
FlinkKafkaConsumerBase.open(Configuration configuration) |
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 |
TwitterSource.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
CollectSink.open(Configuration parameters)
Initialize the connection with the Socket in the server.
|
Modifier and Type | Method and Description |
---|---|
void |
RollingAdditionMapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
static CheckpointedInputGate |
InputProcessorUtil.createCheckpointedInputGate(AbstractInvokable toNotifyOnCheckpoint,
CheckpointingMode checkpointMode,
IOManager ioManager,
InputGate inputGate,
Configuration taskManagerConfig,
String taskName) |
static CheckpointedInputGate[] |
InputProcessorUtil.createCheckpointedInputGatePair(AbstractInvokable toNotifyOnCheckpoint,
CheckpointingMode checkpointMode,
IOManager ioManager,
InputGate inputGate1,
InputGate inputGate2,
Configuration taskManagerConfig,
String taskName) |
Constructor and Description |
---|
StreamOneInputProcessor(InputGate[] inputGates,
TypeSerializer<IN> inputSerializer,
StreamTask<?,?> checkpointedTask,
CheckpointingMode checkpointMode,
Object lock,
IOManager ioManager,
Configuration taskManagerConfig,
StreamStatusMaintainer streamStatusMaintainer,
OneInputStreamOperator<IN,?> streamOperator,
TaskIOMetricGroup metrics,
WatermarkGauge watermarkGauge,
String taskName,
OperatorChain<?,?> operatorChain) |
StreamTwoInputProcessor(Collection<InputGate> inputGates1,
Collection<InputGate> inputGates2,
TypeSerializer<IN1> inputSerializer1,
TypeSerializer<IN2> inputSerializer2,
TwoInputStreamTask<IN1,IN2,?> checkpointedTask,
CheckpointingMode checkpointMode,
Object lock,
IOManager ioManager,
Configuration taskManagerConfig,
StreamStatusMaintainer streamStatusMaintainer,
TwoInputStreamOperator<IN1,IN2,?> streamOperator,
TaskIOMetricGroup metrics,
WatermarkGauge input1WatermarkGauge,
WatermarkGauge input2WatermarkGauge,
String taskName,
OperatorChain<?,?> operatorChain) |
StreamTwoInputSelectableProcessor(Collection<InputGate> inputGates1,
Collection<InputGate> inputGates2,
TypeSerializer<IN1> inputSerializer1,
TypeSerializer<IN2> inputSerializer2,
StreamTask<?,?> streamTask,
CheckpointingMode checkpointingMode,
Object lock,
IOManager ioManager,
Configuration taskManagerConfig,
StreamStatusMaintainer streamStatusMaintainer,
TwoInputStreamOperator<IN1,IN2,?> streamOperator,
WatermarkGauge input1WatermarkGauge,
WatermarkGauge input2WatermarkGauge,
String taskName,
OperatorChain<?,?> operatorChain) |
Modifier and Type | Method and Description |
---|---|
void |
InternalIterableProcessAllWindowFunction.open(Configuration parameters) |
void |
InternalAggregateProcessAllWindowFunction.open(Configuration parameters) |
void |
InternalSingleValueProcessAllWindowFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
BucketingSinkTestProgram.SubtractingMapper.open(Configuration parameters) |
void |
SemanticsCheckMapper.open(Configuration parameters) |
void |
SlidingWindowCheckMapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TableConfig.getConfiguration()
Gives direct access to the underlying key-value map for advanced configuration.
|
Modifier and Type | Method and Description |
---|---|
void |
TableConfig.addConfiguration(Configuration configuration)
Adds the given key-value configuration to the underlying configuration.
|
Constructor and Description |
---|
ExecutionContext(Environment defaultEnvironment,
SessionContext sessionContext,
List<URL> dependencies,
Configuration flinkConfig,
org.apache.commons.cli.Options commandLineOptions,
List<CustomCommandLine<?>> availableCommandLines) |
LocalExecutor(Environment defaultEnvironment,
List<URL> dependencies,
Configuration flinkConfig,
CustomCommandLine<?> commandLine)
Constructor for testing purposes.
|
ResultStore(Configuration flinkConfig) |
Constructor and Description |
---|
BaseHybridHashTable(Configuration conf,
Object owner,
MemoryManager memManager,
long reservedMemorySize,
long preferredMemorySize,
long perRequestMemorySize,
IOManager ioManager,
int avgRecordLen,
long buildRowCount,
boolean tryDistinctBuildRow) |
BinaryHashTable(Configuration conf,
Object owner,
AbstractRowSerializer buildSideSerializer,
AbstractRowSerializer probeSideSerializer,
Projection<BaseRow,BinaryRow> buildSideProjection,
Projection<BaseRow,BinaryRow> probeSideProjection,
MemoryManager memManager,
long reservedMemorySize,
IOManager ioManager,
int avgRecordLen,
int buildRowCount,
boolean useBloomFilters,
HashJoinType type,
JoinCondition condFunc,
boolean reverseJoin,
boolean[] filterNulls,
boolean tryDistinctBuildRow) |
BinaryHashTable(Configuration conf,
Object owner,
AbstractRowSerializer buildSideSerializer,
AbstractRowSerializer probeSideSerializer,
Projection<BaseRow,BinaryRow> buildSideProjection,
Projection<BaseRow,BinaryRow> probeSideProjection,
MemoryManager memManager,
long reservedMemorySize,
long preferredMemorySize,
long perRequestMemorySize,
IOManager ioManager,
int avgRecordLen,
long buildRowCount,
boolean useBloomFilters,
HashJoinType type,
JoinCondition condFunc,
boolean reverseJoin,
boolean[] filterNulls,
boolean tryDistinctBuildRow) |
LongHybridHashTable(Configuration conf,
Object owner,
BinaryRowSerializer buildSideSerializer,
BinaryRowSerializer probeSideSerializer,
MemoryManager memManager,
long reservedMemorySize,
long preferredMemorySize,
long perRequestMemorySize,
IOManager ioManager,
int avgRecordLen,
long buildRowCount) |
Modifier and Type | Method and Description |
---|---|
void |
GroupAggFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
DeduplicateKeepFirstRowFunction.open(Configuration configure) |
void |
DeduplicateKeepLastRowFunction.open(Configuration configure) |
Modifier and Type | Method and Description |
---|---|
TableFunctionResultFuture<BaseRow> |
AsyncLookupJoinWithCalcRunner.createFetcherResultFuture(Configuration parameters) |
TableFunctionResultFuture<BaseRow> |
AsyncLookupJoinRunner.createFetcherResultFuture(Configuration parameters) |
void |
AsyncLookupJoinWithCalcRunner.open(Configuration parameters) |
void |
LookupJoinRunner.open(Configuration parameters) |
void |
AsyncLookupJoinRunner.open(Configuration parameters) |
void |
LookupJoinWithCalcRunner.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
IterativeConditionRunner.open(Configuration parameters) |
void |
PatternProcessFunctionRunner.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractRowTimeUnboundedPrecedingOver.open(Configuration parameters) |
void |
ProcTimeRangeBoundedPrecedingFunction.open(Configuration parameters) |
void |
ProcTimeRowsBoundedPrecedingFunction.open(Configuration parameters) |
void |
RowTimeRangeBoundedPrecedingFunction.open(Configuration parameters) |
void |
RowTimeRowsBoundedPrecedingFunction.open(Configuration parameters) |
void |
ProcTimeUnboundedPrecedingFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
AppendOnlyTopNFunction.open(Configuration parameters) |
void |
UpdatableTopNFunction.open(Configuration parameters) |
void |
AbstractTopNFunction.open(Configuration parameters) |
void |
RetractableTopNFunction.open(Configuration parameters) |
Constructor and Description |
---|
BinaryExternalSorter(Object owner,
MemoryManager memoryManager,
long reservedMemorySize,
IOManager ioManager,
AbstractRowSerializer<BaseRow> inputSerializer,
BinaryRowSerializer serializer,
NormalizedKeyComputer normalizedKeyComputer,
RecordComparator comparator,
Configuration conf) |
BinaryExternalSorter(Object owner,
MemoryManager memoryManager,
long reservedMemorySize,
IOManager ioManager,
AbstractRowSerializer<BaseRow> inputSerializer,
BinaryRowSerializer serializer,
NormalizedKeyComputer normalizedKeyComputer,
RecordComparator comparator,
Configuration conf,
float startSpillingFraction) |
BufferedKVExternalSorter(IOManager ioManager,
BinaryRowSerializer keySerializer,
BinaryRowSerializer valueSerializer,
NormalizedKeyComputer nKeyComputer,
RecordComparator comparator,
int pageSize,
Configuration conf) |
Modifier and Type | Method and Description |
---|---|
void |
StatefulStreamingJob.MyStatefulFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
SecureTestEnvironment.populateFlinkSecureConfigurations(Configuration flinkConf) |
Modifier and Type | Method and Description |
---|---|
TestProcessBuilder |
TestProcessBuilder.addConfigAsMainClassArgs(Configuration config) |
static Configuration |
SecureTestEnvironment.populateFlinkSecureConfigurations(Configuration flinkConf) |
MiniClusterResourceConfiguration.Builder |
MiniClusterResourceConfiguration.Builder.setConfiguration(Configuration configuration) |
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(Configuration... testConfigs) |
Modifier and Type | Method and Description |
---|---|
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(List<Configuration> testConfigs) |
Modifier and Type | Method and Description |
---|---|
void |
FlinkDistribution.appendConfiguration(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 |
AbstractYarnClusterDescriptor.getFlinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
static int |
Utils.calculateHeapSize(int memory,
Configuration conf)
See documentation.
|
protected ClusterClient<org.apache.hadoop.yarn.api.records.ApplicationId> |
YarnClusterDescriptor.createYarnClusterClient(AbstractYarnClusterDescriptor descriptor,
int numberTaskManagers,
int slotsPerTaskManager,
org.apache.hadoop.yarn.api.records.ApplicationReport report,
Configuration flinkConfiguration,
boolean perJobCluster) |
protected abstract ClusterClient<org.apache.hadoop.yarn.api.records.ApplicationId> |
AbstractYarnClusterDescriptor.createYarnClusterClient(AbstractYarnClusterDescriptor descriptor,
int numberTaskManagers,
int slotsPerTaskManager,
org.apache.hadoop.yarn.api.records.ApplicationReport report,
Configuration flinkConfiguration,
boolean perJobCluster)
Creates a YarnClusterClient; may be overridden 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.
|
org.apache.hadoop.yarn.api.records.ApplicationReport |
AbstractYarnClusterDescriptor.startAppMaster(Configuration configuration,
String applicationName,
String yarnClusterEntrypoint,
JobGraph jobGraph,
org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
org.apache.hadoop.yarn.client.api.YarnClientApplication yarnApplication,
ClusterSpecification clusterSpecification) |
Constructor and Description |
---|
AbstractYarnClusterDescriptor(Configuration flinkConfiguration,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfiguration,
String configurationDirectory,
org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
boolean sharedYarnClient) |
YarnClusterDescriptor(Configuration flinkConfiguration,
org.apache.hadoop.yarn.conf.YarnConfiguration yarnConfiguration,
String configurationDirectory,
org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
boolean sharedYarnClient) |
YarnResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
Configuration flinkConfig,
Map<String,String> env,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
protected Configuration |
FlinkYarnSessionCli.applyCommandLineOptionsToConfiguration(org.apache.commons.cli.CommandLine commandLine) |
Constructor and Description |
---|
FlinkYarnSessionCli(Configuration configuration,
String configurationDirectory,
String shortPrefix,
String longPrefix) |
FlinkYarnSessionCli(Configuration configuration,
String configurationDirectory,
String shortPrefix,
String longPrefix,
boolean acceptInteractiveInput) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
YarnEntrypointUtils.loadConfiguration(String workingDirectory,
Map<String,String> env,
org.slf4j.Logger log) |
Modifier and Type | Method and Description |
---|---|
ResourceManager<YarnWorkerNode> |
YarnResourceManagerFactory.createActiveResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
protected DispatcherResourceManagerComponentFactory<?> |
YarnJobClusterEntrypoint.createDispatcherResourceManagerComponentFactory(Configuration configuration) |
protected DispatcherResourceManagerComponentFactory<?> |
YarnSessionClusterEntrypoint.createDispatcherResourceManagerComponentFactory(Configuration configuration) |
protected String |
YarnJobClusterEntrypoint.getRPCPortRange(Configuration configuration) |
protected String |
YarnSessionClusterEntrypoint.getRPCPortRange(Configuration configuration) |
protected SecurityContext |
YarnJobClusterEntrypoint.installSecurityContext(Configuration configuration) |
protected SecurityContext |
YarnSessionClusterEntrypoint.installSecurityContext(Configuration configuration) |
static SecurityContext |
YarnEntrypointUtils.installSecurityContext(Configuration configuration,
String workingDirectory) |
Constructor and Description |
---|
YarnJobClusterEntrypoint(Configuration configuration,
String workingDirectory) |
YarnSessionClusterEntrypoint(Configuration configuration,
String workingDirectory) |
Modifier and Type | Method and Description |
---|---|
static YarnHighAvailabilityServices |
YarnHighAvailabilityServices.forSingleJobAppMaster(Configuration flinkConfig,
Configuration hadoopConfig)
Creates the high-availability services for a single-job Flink YARN application, to be
used in the Application Master that runs both ResourceManager and JobManager.
|
static YarnHighAvailabilityServices |
YarnHighAvailabilityServices.forYarnTaskManager(Configuration flinkConfig,
Configuration hadoopConfig)
Creates the high-availability services for the TaskManagers participating in
a Flink YARN application.
|
Constructor and Description |
---|
AbstractYarnNonHaServices(Configuration config,
Configuration hadoopConf)
Creates new YARN high-availability services, configuring the file system and recovery
data directory based on the working directory in the given Hadoop configuration.
|
YarnHighAvailabilityServices(Configuration config,
Configuration hadoopConf)
Creates new YARN high-availability services, configuring the file system and recovery
data directory based on the working directory in the given Hadoop configuration.
|
YarnIntraNonHaMasterServices(Configuration config,
Configuration hadoopConf)
Creates new YarnIntraNonHaMasterServices for the given Flink and YARN configuration.
|
YarnPreConfiguredMasterNonHaServices(Configuration config,
Configuration hadoopConf,
HighAvailabilityServicesUtils.AddressResolution addressResolution)
Creates new YarnPreConfiguredMasterHaServices for the given Flink and YARN configuration.
|
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.