Modifier and Type | Method and Description |
---|---|
void |
TableInputFormat.configure(Configuration parameters)
creates a
Scan object and a HTable connection |
Modifier and Type | Method and Description |
---|---|
static PlanExecutor |
PlanExecutor.createLocalExecutor(Configuration configuration)
Creates an executor that runs the plan locally in a multi-threaded environment.
|
static PlanExecutor |
PlanExecutor.createRemoteExecutor(String hostname,
int port,
Configuration clientConfiguration,
URL[] jarFiles,
URL[] globalClasspaths)
Creates an executor that runs the plan on a remote environment.
|
Modifier and Type | Method and Description |
---|---|
static Set<Map.Entry<String,DistributedCache.DistributedCacheEntry>> |
DistributedCache.readFileInfoFromConfig(Configuration conf) |
static void |
DistributedCache.writeFileInfoToConfig(String name,
DistributedCache.DistributedCacheEntry e,
Configuration conf) |
Modifier and Type | Method and Description |
---|---|
void |
RichFunction.open(Configuration parameters)
Initialization method for the function.
|
void |
AbstractRichFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
static void |
FunctionUtils.openFunction(Function function,
Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
ReplicatingInputFormat.configure(Configuration parameters) |
void |
OutputFormat.configure(Configuration parameters)
Configures this output format.
|
void |
InputFormat.configure(Configuration parameters)
Configures this input format.
|
void |
GenericInputFormat.configure(Configuration parameters) |
void |
FileOutputFormat.configure(Configuration parameters) |
void |
FileInputFormat.configure(Configuration parameters)
Configures the file input format by reading the file path from the configuration.
|
void |
DelimitedInputFormat.configure(Configuration parameters)
Configures this input format by reading the path to the file from the configuration andge the string that
defines the record delimiter.
|
void |
BinaryOutputFormat.configure(Configuration parameters) |
void |
BinaryInputFormat.configure(Configuration parameters) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
Operator.parameters |
Modifier and Type | Method and Description |
---|---|
Configuration |
Operator.getParameters()
Gets the stub parameters of this contract.
|
Modifier and Type | Method and Description |
---|---|
void |
BulkIterationBase.TerminationCriterionMapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
TypeSerializerFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
TypeComparatorFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
TypeSerializerFactory.writeParametersToConfig(Configuration config) |
void |
TypeComparatorFactory.writeParametersToConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
Utils.CountHelper.configure(Configuration parameters) |
void |
Utils.CollectHelper.configure(Configuration parameters) |
void |
Utils.ChecksumHashCodeHelper.configure(Configuration parameters) |
static LocalEnvironment |
ExecutionEnvironment.createLocalEnvironment(Configuration customConfiguration)
Creates a
LocalEnvironment . |
static ExecutionEnvironment |
ExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
String... jarFiles)
Creates a
RemoteEnvironment . |
Constructor and Description |
---|
LocalEnvironment(Configuration config)
Creates a new local environment that configures its local executor with the given configuration.
|
RemoteEnvironment(String host,
int port,
Configuration clientConfig,
String[] jarFiles)
Creates a new RemoteEnvironment that points to the master (JobManager) described by the
given host name and port.
|
RemoteEnvironment(String host,
int port,
Configuration clientConfig,
String[] jarFiles,
URL[] globalClasspaths)
Creates a new RemoteEnvironment that points to the master (JobManager) described by the
given host name and port.
|
Modifier and Type | Method and Description |
---|---|
void |
HadoopOutputFormatBase.configure(Configuration parameters) |
void |
HadoopInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopOutputFormatBase.configure(Configuration parameters) |
void |
HadoopInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
TextValueInputFormat.configure(Configuration parameters) |
void |
TextInputFormat.configure(Configuration parameters) |
void |
PrintingOutputFormat.configure(Configuration parameters) |
void |
LocalCollectionOutputFormat.configure(Configuration parameters) |
void |
DiscardingOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
JDBCInputFormat.configure(Configuration parameters) |
void |
JDBCOutputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
Configuration |
UdfOperator.getParameters()
Gets the configuration parameters that will be passed to the UDF's open method
AbstractRichFunction.open(Configuration) . |
Configuration |
TwoInputUdfOperator.getParameters() |
Configuration |
SingleInputUdfOperator.getParameters() |
Configuration |
DataSource.getParameters() |
Configuration |
DataSink.getParameters() |
Modifier and Type | Method and Description |
---|---|
void |
AggregateOperator.AggregatingUdf.open(Configuration parameters) |
O |
UdfOperator.withParameters(Configuration parameters)
Sets the configuration parameters for the UDF.
|
O |
TwoInputUdfOperator.withParameters(Configuration parameters) |
O |
SingleInputUdfOperator.withParameters(Configuration parameters) |
DataSource<OUT> |
DataSource.withParameters(Configuration parameters)
Pass a configuration to the InputFormat
|
DataSink<T> |
DataSink.withParameters(Configuration parameters)
Pass a configuration to the OutputFormat
|
Modifier and Type | Method and Description |
---|---|
void |
WrappingFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
RuntimeSerializerFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeComparatorFactory.readParametersFromConfig(Configuration config,
ClassLoader cl) |
void |
RuntimeSerializerFactory.writeParametersToConfig(Configuration config) |
void |
RuntimeComparatorFactory.writeParametersToConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ParameterTool.getConfiguration()
Returns a
Configuration object from this ParameterTool |
Modifier and Type | Method and Description |
---|---|
static ExecutionEnvironment |
ExecutionEnvironment.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
ExecutionEnvironment |
ExecutionEnvironment$.createLocalEnvironment(Configuration customConfiguration)
Creates a local execution environment.
|
static ExecutionEnvironment |
ExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
scala.collection.Seq<String> jarFiles)
Creates a remote execution environment.
|
ExecutionEnvironment |
ExecutionEnvironment$.createRemoteEnvironment(String host,
int port,
Configuration clientConfiguration,
scala.collection.Seq<String> jarFiles)
Creates a remote execution environment.
|
DataSet<T> |
DataSet.withParameters(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
ScalaAggregateOperator.AggregatingUdf.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
ExpressionFilterFunction.open(Configuration c) |
void |
ExpressionSelectFunction.open(Configuration c) |
void |
ExpressionAggregateFunction.open(Configuration conf) |
void |
ExpressionJoinFunction.open(Configuration c) |
Modifier and Type | Method and Description |
---|---|
Configuration |
CliFrontend.getConfiguration()
Getter which returns a copy of the associated configuration
|
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) |
Constructor and Description |
---|
Client(Configuration config)
Creates a instance that submits the programs to the JobManager defined in the
configuration.
|
Client(Configuration config,
int maxSlots)
Creates a new instance of the class that submits the jobs to a job-manager.
|
Modifier and Type | Class and Description |
---|---|
class |
UnmodifiableConfiguration
Unmodifiable version of the Configuration class.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
Configuration.clone() |
static Configuration |
GlobalConfiguration.getConfiguration()
Gets a
Configuration object with the values of this GlobalConfiguration |
Modifier and Type | Method and Description |
---|---|
void |
UnmodifiableConfiguration.addAll(Configuration other) |
void |
Configuration.addAll(Configuration other) |
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 void |
GlobalConfiguration.includeConfiguration(Configuration conf)
Merges the given
Configuration object into the global
configuration. |
Constructor and Description |
---|
Configuration(Configuration other)
Creates a new configuration with the copy of the given configuration.
|
UnmodifiableConfiguration(Configuration config)
Creates a new UnmodifiableConfiguration, which holds a copy of the given configuration
that cannot be altered.
|
Modifier and Type | Method and Description |
---|---|
static void |
FileSystem.setDefaultScheme(Configuration config)
Sets the default filesystem scheme based on the user-specified configuration parameter
fs.default-scheme . |
Modifier and Type | Method and Description |
---|---|
void |
KMeans.SelectNearestCenter.open(Configuration parameters)
Reads the centroid values from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
FileCopyTaskInputFormat.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
LinearRegression.SubUpdate.open(Configuration parameters)
Reads the parameters from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
EmptyFieldsCountAccumulator.EmptyFieldFilter.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
KMeans.SelectNearestCenter.open(Configuration parameters)
Reads the centroid values from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
LinearRegression.SubUpdate.open(Configuration parameters)
Reads the parameters from a broadcast variable into a collection.
|
Modifier and Type | Method and Description |
---|---|
void |
ExampleUtils.PrintingOutputFormatWithMessage.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopReduceFunction.open(Configuration parameters) |
void |
HadoopMapFunction.open(Configuration parameters) |
void |
HadoopReduceCombineFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HCatInputFormatBase.configure(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
HCatInputFormat.configure(Configuration parameters) |
Constructor and Description |
---|
Optimizer(Configuration config)
Creates a new optimizer instance.
|
Optimizer(CostEstimator estimator,
Configuration config)
Creates a new optimizer instance.
|
Optimizer(DataStatistics stats,
Configuration config)
Creates a new optimizer instance that uses the statistics object to determine properties about the input.
|
Optimizer(DataStatistics stats,
CostEstimator estimator,
Configuration config)
Creates a new optimizer instance that uses the statistics object to determine properties about the input.
|
Constructor and Description |
---|
JobGraphGenerator(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
PythonMapPartition.open(Configuration config)
Opens this function.
|
void |
PythonCoGroup.open(Configuration config)
Opens this function.
|
Modifier and Type | Method and Description |
---|---|
void |
PythonStreamer.sendBroadCastVariables(Configuration config)
Sends all broadcast-variables encoded in the configuration to the external process.
|
Modifier and Type | Method and Description |
---|---|
akka.actor.ActorSystem |
AkkaUtils$.createActorSystem(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an actor system.
|
static akka.actor.ActorSystem |
AkkaUtils.createActorSystem(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an actor system.
|
akka.actor.ActorSystem |
AkkaUtils$.createLocalActorSystem(Configuration configuration)
Creates a local actor system without remoting.
|
static akka.actor.ActorSystem |
AkkaUtils.createLocalActorSystem(Configuration configuration)
Creates a local actor system without remoting.
|
com.typesafe.config.Config |
AkkaUtils$.getAkkaConfig(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an akka config with the provided configuration values.
|
static com.typesafe.config.Config |
AkkaUtils.getAkkaConfig(Configuration configuration,
scala.Option<scala.Tuple2<String,Object>> listeningAddress)
Creates an akka config with the provided configuration values.
|
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getClientTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getClientTimeout(Configuration config) |
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getLookupTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getLookupTimeout(Configuration config) |
scala.concurrent.duration.FiniteDuration |
AkkaUtils$.getTimeout(Configuration config) |
static scala.concurrent.duration.FiniteDuration |
AkkaUtils.getTimeout(Configuration config) |
Constructor and Description |
---|
BlobCache(InetSocketAddress serverAddress,
Configuration configuration) |
BlobServer(Configuration config)
Instantiates a new BLOB server and binds it to a free network port.
|
Modifier and Type | Method and Description |
---|---|
static SavepointStore |
SavepointStoreFactory.createFromConfig(Configuration config)
Creates a
SavepointStore from the specified Configuration. |
Constructor and Description |
---|
ZooKeeperCheckpointRecoveryFactory(org.apache.curator.framework.CuratorFramework client,
Configuration config) |
Modifier and Type | Method and Description |
---|---|
static akka.actor.ActorSystem |
JobClient.startJobClientActorSystem(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskDeploymentDescriptor.getJobConfiguration()
Returns the configuration of the job the task belongs to.
|
Configuration |
TaskDeploymentDescriptor.getTaskConfiguration()
Returns the task's configuration object.
|
Constructor and Description |
---|
TaskDeploymentDescriptor(JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionId,
String taskName,
int indexInSubtaskGroup,
int numberOfSubtasks,
int attemptNumber,
Configuration jobConfiguration,
Configuration taskConfiguration,
String invokableClassName,
List<ResultPartitionDeploymentDescriptor> producedPartitions,
List<InputGateDeploymentDescriptor> inputGates,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
int targetSlotNumber) |
TaskDeploymentDescriptor(JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionId,
String taskName,
int indexInSubtaskGroup,
int numberOfSubtasks,
int attemptNumber,
Configuration jobConfiguration,
Configuration taskConfiguration,
String invokableClassName,
List<ResultPartitionDeploymentDescriptor> producedPartitions,
List<InputGateDeploymentDescriptor> inputGates,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
long recoveryTimestamp)
Constructs a task deployment descriptor.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
Environment.getJobConfiguration()
Returns the job-wide configuration object that was attached to the JobGraph.
|
Configuration |
Environment.getTaskConfiguration()
Returns the task-wide configuration object, originally attache to the job vertex.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
ExecutionGraph.getJobConfiguration() |
Constructor and Description |
---|
ExecutionGraph(scala.concurrent.ExecutionContext executionContext,
JobID jobId,
String jobName,
Configuration jobConfig,
scala.concurrent.duration.FiniteDuration timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
ClassLoader userClassLoader) |
Modifier and Type | Method and Description |
---|---|
static NoRestartStrategy |
NoRestartStrategy.create(Configuration configuration)
Creates a NoRestartStrategy instance.
|
static FixedDelayRestartStrategy |
FixedDelayRestartStrategy.create(Configuration configuration)
Creates a FixedDelayRestartStrategy from the given Configuration.
|
static RestartStrategy |
RestartStrategyFactory.createFromConfig(Configuration configuration)
Creates a
RestartStrategy instance from the given Configuration . |
Constructor and Description |
---|
FileCache(Configuration config) |
Constructor and Description |
---|
NettyConfig(InetAddress serverAddress,
int serverPort,
int memorySegmentSize,
int numberOfSlots,
Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
JobVertex.getConfiguration()
Returns the vertex's configuration object which can be used to pass custom settings to the task at runtime.
|
Configuration |
JobGraph.getJobConfiguration()
Returns the configuration object for this job.
|
Modifier and Type | Method and Description |
---|---|
Configuration |
AbstractInvokable.getJobConfiguration()
Returns the job configuration object which was attached to the original
JobGraph . |
Configuration |
AbstractInvokable.getTaskConfiguration()
Returns the task configuration object which was attached to the original
JobVertex . |
Modifier and Type | Method and Description |
---|---|
protected Configuration |
JobManager.flinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
static scala.Tuple4<Configuration,JobManagerMode,String,Iterator<Integer>> |
JobManager.parseArgs(String[] args)
Loads the configuration, execution mode and the listening address from the provided command
line arguments.
|
scala.Tuple4<Configuration,JobManagerMode,String,Iterator<Integer>> |
JobManager$.parseArgs(String[] args)
Loads the configuration, execution mode and the listening address from the provided command
line arguments.
|
Modifier and Type | Method and Description |
---|---|
static scala.Tuple12<ExecutorService,InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategy,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,SavepointStore,scala.concurrent.duration.FiniteDuration> |
JobManager.createJobManagerComponents(Configuration configuration,
scala.Option<LeaderElectionService> leaderElectionServiceOption)
Create the job manager components as (instanceManager, scheduler, libraryCacheManager,
archiverProps, defaultExecutionRetries,
delayBetweenRetries, timeout)
|
scala.Tuple12<ExecutorService,InstanceManager,Scheduler,BlobLibraryCacheManager,RestartStrategy,scala.concurrent.duration.FiniteDuration,Object,LeaderElectionService,SubmittedJobGraphStore,CheckpointRecoveryFactory,SavepointStore,scala.concurrent.duration.FiniteDuration> |
JobManager$.createJobManagerComponents(Configuration configuration,
scala.Option<LeaderElectionService> leaderElectionServiceOption)
Create the job manager components as (instanceManager, scheduler, libraryCacheManager,
archiverProps, defaultExecutionRetries,
delayBetweenRetries, timeout)
|
static RecoveryMode |
RecoveryMode.fromConfig(Configuration config)
Return the configured
RecoveryMode . |
static akka.actor.ActorRef |
JobManager.getJobManagerActorRef(InetSocketAddress address,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
akka.actor.ActorRef |
JobManager$.getJobManagerActorRef(InetSocketAddress address,
akka.actor.ActorSystem system,
Configuration config)
Resolves the JobManager actor reference in a blocking fashion.
|
static String |
JobManager.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
String |
JobManager$.getRemoteJobManagerAkkaURL(Configuration config)
Returns the JobManager actor's remote Akka URL, given the configured hostname and port.
|
static boolean |
RecoveryMode.isHighAvailabilityModeActivated(Configuration configuration)
Returns true if the defined recovery mode supports high availability.
|
static void |
JobManager.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort)
Starts and runs the JobManager with all its components.
|
void |
JobManager$.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort)
Starts and runs the JobManager with all its components.
|
static void |
JobManager.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
Iterator<Integer> listeningPortRange)
Starts and runs the JobManager with all its components trying to bind to
a port in the specified range.
|
void |
JobManager$.runJobManager(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
Iterator<Integer> listeningPortRange)
Starts and runs the JobManager with all its components trying to bind to
a port in the specified range.
|
static scala.Tuple4<akka.actor.ActorSystem,akka.actor.ActorRef,akka.actor.ActorRef,scala.Option<WebMonitor>> |
JobManager.startActorSystemAndJobManagerActors(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts an ActorSystem, the JobManager and all its components including the WebMonitor.
|
scala.Tuple4<akka.actor.ActorSystem,akka.actor.ActorRef,akka.actor.ActorRef,scala.Option<WebMonitor>> |
JobManager$.startActorSystemAndJobManagerActors(Configuration configuration,
JobManagerMode executionMode,
String listeningAddress,
int listeningPort,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts an ActorSystem, the JobManager and all its components including the WebMonitor.
|
static scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager$.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
static scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
scala.Option<String> jobMangerActorName,
scala.Option<String> archiveActorName,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
scala.Tuple2<akka.actor.ActorRef,akka.actor.ActorRef> |
JobManager$.startJobManagerActors(Configuration configuration,
akka.actor.ActorSystem actorSystem,
scala.Option<String> jobMangerActorName,
scala.Option<String> archiveActorName,
Class<? extends JobManager> jobManagerClass,
Class<? extends MemoryArchivist> archiveClass)
Starts the JobManager and job archiver based on the given configuration, in the
given actor system.
|
Constructor and Description |
---|
JobManager(Configuration flinkConfiguration,
ExecutorService executorService,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategy defaultRestartStrategy,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout) |
Modifier and Type | Method and Description |
---|---|
Configuration |
FlinkMiniCluster.configuration() |
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Configuration |
LocalFlinkMiniCluster.getDefaultConfig() |
Configuration |
FlinkMiniCluster.userConfiguration() |
Modifier and Type | Method and Description |
---|---|
Configuration |
LocalFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
abstract Configuration |
FlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
void |
LocalFlinkMiniCluster.initializeIOFormatClasses(Configuration configuration) |
void |
FlinkMiniCluster.setDefaultCiConfig(Configuration config)
Sets CI environment (Travis) specific config defaults.
|
void |
LocalFlinkMiniCluster.setMemory(Configuration config) |
scala.Option<WebMonitor> |
FlinkMiniCluster.startWebServer(Configuration config,
akka.actor.ActorSystem actorSystem,
String jobManagerAkkaURL) |
Constructor and Description |
---|
FlinkMiniCluster(Configuration userConfiguration,
boolean useSingleActorSystem) |
LocalFlinkMiniCluster(Configuration userConfiguration) |
LocalFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
Modifier and Type | Method and Description |
---|---|
static void |
BatchTask.openUserCode(Function stub,
Configuration parameters)
Opens the given stub using its
RichFunction.open(Configuration) method. |
Modifier and Type | Method and Description |
---|---|
void |
CombiningUnilateralSortMerger.setUdfConfiguration(Configuration config) |
Modifier and Type | Class and Description |
---|---|
static class |
TaskConfig.DelegatingConfiguration
A configuration that manages a subset of keys with a common prefix from a given configuration.
|
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.DelegatingConfiguration.addAll(Configuration other) |
void |
TaskConfig.DelegatingConfiguration.addAll(Configuration other,
String prefix) |
void |
TaskConfig.setStubParameters(Configuration parameters) |
Constructor and Description |
---|
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.
|
TaskConfig(Configuration config)
Creates a new Task Config that wraps the given configuration.
|
Modifier and Type | Method and Description |
---|---|
AbstractStateBackend |
StateBackendFactory.createFromConfig(Configuration config)
Creates the state backend, optionally using the given configuration.
|
Modifier and Type | Method and Description |
---|---|
FsStateBackend |
FsStateBackendFactory.createFromConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
TaskManagerConfiguration.configuration() |
Configuration |
TaskManagerRuntimeInfo.getConfiguration()
Gets the configuration that the TaskManager was started with.
|
Configuration |
Task.getJobConfiguration() |
Configuration |
RuntimeEnvironment.getJobConfiguration() |
Configuration |
Task.getTaskConfiguration() |
Configuration |
RuntimeEnvironment.getTaskConfiguration() |
Configuration |
TaskManager$.parseArgsAndLoadConfig(String[] args)
Parse the command line arguments of the TaskManager and loads the configuration.
|
static Configuration |
TaskManager.parseArgsAndLoadConfig(String[] args)
Parse the command line arguments of the TaskManager and loads the configuration.
|
Modifier and Type | Method and Description |
---|---|
scala.Tuple2<String,Object> |
TaskManager$.getAndCheckJobManagerAddress(Configuration configuration)
Gets the hostname and port of the JobManager from the configuration.
|
static scala.Tuple2<String,Object> |
TaskManager.getAndCheckJobManagerAddress(Configuration configuration)
Gets the hostname and port of the JobManager from the configuration.
|
scala.Tuple4<TaskManagerConfiguration,NetworkEnvironmentConfiguration,InstanceConnectionInfo,MemoryType> |
TaskManager$.parseTaskManagerConfiguration(Configuration configuration,
String taskManagerHostname,
boolean localTaskManagerCommunication)
Utility method to extract TaskManager config parameters from the configuration and to
sanity check them.
|
static scala.Tuple4<TaskManagerConfiguration,NetworkEnvironmentConfiguration,InstanceConnectionInfo,MemoryType> |
TaskManager.parseTaskManagerConfiguration(Configuration configuration,
String taskManagerHostname,
boolean localTaskManagerCommunication)
Utility method to extract TaskManager config parameters from the configuration and to
sanity check them.
|
void |
TaskManager$.runTaskManager(String taskManagerHostname,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
void |
TaskManager$.runTaskManager(String taskManagerHostname,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
scala.Tuple2<String,Object> |
TaskManager$.selectNetworkInterfaceAndPort(Configuration configuration) |
static scala.Tuple2<String,Object> |
TaskManager.selectNetworkInterfaceAndPort(Configuration configuration) |
void |
TaskManager$.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
akka.actor.ActorRef |
TaskManager$.startTaskManagerComponentsAndActor(Configuration configuration,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass) |
static akka.actor.ActorRef |
TaskManager.startTaskManagerComponentsAndActor(Configuration configuration,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass) |
Constructor and Description |
---|
RuntimeEnvironment(JobID jobId,
JobVertexID jobVertexId,
ExecutionAttemptID executionId,
TaskInfo taskInfo,
Configuration jobConfiguration,
Configuration taskConfiguration,
ClassLoader userCodeClassLoader,
MemoryManager memManager,
IOManager ioManager,
BroadcastVariableManager bcVarManager,
AccumulatorRegistry accumulatorRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
ActorGateway jobManager,
TaskManagerRuntimeInfo taskManagerInfo) |
TaskManagerConfiguration(String[] tmpDirPaths,
long cleanupInterval,
scala.concurrent.duration.FiniteDuration timeout,
scala.Option<scala.concurrent.duration.FiniteDuration> maxRegistrationDuration,
int numberOfSlots,
Configuration configuration) |
TaskManagerRuntimeInfo(String hostname,
Configuration configuration)
Creates a runtime info.
|
Modifier and Type | Method and Description |
---|---|
static ZooKeeperCheckpointIDCounter |
ZooKeeperUtils.createCheckpointIDCounter(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId)
Creates a
ZooKeeperCheckpointIDCounter instance. |
static CompletedCheckpointStore |
ZooKeeperUtils.createCompletedCheckpoints(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId,
int maxNumberOfCheckpointsToRetain,
ClassLoader userClassLoader)
Creates a
ZooKeeperCompletedCheckpointStore instance. |
static ZooKeeperLeaderElectionService |
ZooKeeperUtils.createLeaderElectionService(Configuration configuration)
Creates a
ZooKeeperLeaderElectionService instance and a new CuratorFramework client. |
static ZooKeeperLeaderElectionService |
ZooKeeperUtils.createLeaderElectionService(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperLeaderElectionService instance. |
static ZooKeeperLeaderRetrievalService |
ZooKeeperUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
ZooKeeperLeaderRetrievalService instance. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
StandaloneLeaderRetrievalService from the given configuration. |
static LeaderRetrievalService |
LeaderRetrievalUtils.createLeaderRetrievalService(Configuration configuration)
Creates a
LeaderRetrievalService based on the provided Configuration object. |
static StandaloneLeaderRetrievalService |
StandaloneUtils.createLeaderRetrievalService(Configuration configuration,
String jobManagerName)
Creates a
StandaloneLeaderRetrievalService form the given configuration and the
JobManager name. |
static ZooKeeperSubmittedJobGraphStore |
ZooKeeperUtils.createSubmittedJobGraphs(org.apache.curator.framework.CuratorFramework client,
Configuration configuration)
Creates a
ZooKeeperSubmittedJobGraphStore instance. |
static String |
ZooKeeperUtils.getZooKeeperEnsemble(Configuration flinkConf)
Returns the configured ZooKeeper quorum (and removes whitespace, because ZooKeeper does not
tolerate it).
|
static boolean |
ZooKeeperUtils.isZooKeeperRecoveryMode(Configuration flinkConf)
Returns whether
RecoveryMode.ZOOKEEPER is configured. |
static org.apache.curator.framework.CuratorFramework |
ZooKeeperUtils.startCuratorFramework(Configuration configuration)
Starts a
CuratorFramework instance and connects it to the given ZooKeeper
quorum. |
Modifier and Type | Method and Description |
---|---|
static WebMonitorUtils.LogFileLocation |
WebMonitorUtils.LogFileLocation.find(Configuration config)
Finds the Flink log directory using log.file Java property that is set during startup.
|
static WebMonitor |
WebMonitorUtils.startWebRuntimeMonitor(Configuration config,
LeaderRetrievalService leaderRetrievalService,
akka.actor.ActorSystem actorSystem)
Starts the web runtime monitor.
|
Constructor and Description |
---|
WebMonitorConfig(Configuration config) |
WebRuntimeMonitor(Configuration config,
LeaderRetrievalService leaderRetrievalService,
akka.actor.ActorSystem actorSystem) |
Constructor and Description |
---|
JobManagerConfigHandler(Configuration config) |
Modifier and Type | Method and Description |
---|---|
abstract Configuration |
AbstractFlinkYarnCluster.getFlinkConfiguration()
Return the Flink configuration object
|
abstract Configuration |
AbstractFlinkYarnClient.getFlinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
abstract void |
AbstractFlinkYarnClient.setFlinkConfiguration(Configuration conf)
Flink configuration
|
Modifier and Type | Method and Description |
---|---|
static LocalStreamEnvironment |
StreamExecutionEnvironment.createLocalEnvironment(int parallelism,
Configuration configuration)
Creates a
LocalStreamEnvironment . |
static StreamExecutionEnvironment |
StreamExecutionEnvironment.createRemoteEnvironment(String host,
int port,
Configuration clientConfig,
String... jarFiles)
Creates a
RemoteStreamEnvironment . |
Constructor and Description |
---|
LocalStreamEnvironment(Configuration config)
Creates a new local stream environment that configures its local executor with the given configuration.
|
RemoteStreamEnvironment(String host,
int port,
Configuration config,
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 config,
String[] jarFiles,
URL[] globalClasspaths)
Creates a new RemoteStreamEnvironment that points to the master
(JobManager) described by the given host name and port.
|
Modifier and Type | Method and Description |
---|---|
void |
SocketClientSink.open(Configuration parameters)
Initialize the connection with the Socket in the server.
|
void |
PrintSinkFunction.open(Configuration parameters) |
void |
OutputFormatSinkFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
MultipleIdsMessageAcknowledgingSourceBase.open(Configuration parameters) |
void |
MessageAcknowledgingSourceBase.open(Configuration parameters) |
void |
FromSplittableIteratorFunction.open(Configuration parameters) |
void |
FileSourceFunction.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
FoldApplyWindowFunction.open(Configuration configuration) |
void |
FoldApplyAllWindowFunction.open(Configuration configuration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
StreamConfig.getConfiguration() |
Constructor and Description |
---|
StreamConfig(Configuration config) |
Modifier and Type | Method and Description |
---|---|
Configuration |
AbstractUdfStreamOperator.getUserFunctionParameters()
Since the streaming API does not implement any parametrization of functions via a
configuration, the config returned here is actually empty.
|
Modifier and Type | Method and Description |
---|---|
void |
StatefulFunction.open(Configuration c) |
Modifier and Type | Method and Description |
---|---|
void |
ScalaWindowFunctionWrapper.open(Configuration parameters) |
void |
ScalaAllWindowFunctionWrapper.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
ElasticsearchSink.open(Configuration configuration)
Initializes the connection to Elasticsearch by either creating an embedded
Node and retrieving the
Client from it or by creating a
TransportClient . |
Modifier and Type | Method and Description |
---|---|
void |
FlumeSink.open(Configuration config) |
Modifier and Type | Method and Description |
---|---|
void |
RollingSink.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
FlinkKafkaConsumer09.open(Configuration parameters) |
void |
FlinkKafkaProducerBase.open(Configuration configuration)
Initializes the connection to Kafka.
|
void |
FlinkKafkaConsumer08.open(Configuration parameters) |
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 |
WikipediaEditsSource.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
WindowJoin.SalarySource.open(Configuration parameters) |
Modifier and Type | Method and Description |
---|---|
void |
InternalSingleValueWindowFunction.open(Configuration parameters) |
void |
InternalIterableWindowFunction.open(Configuration parameters) |
Modifier and Type | Field and Description |
---|---|
protected Configuration |
AbstractTestBase.config
Configuration to start the testing cluster with
|
Modifier and Type | Method and Description |
---|---|
Configuration |
ForkableFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ForkableFlinkMiniCluster.generateConfiguration(Configuration userConfiguration) |
static ForkableFlinkMiniCluster |
TestBaseUtils.startCluster(Configuration config,
boolean singleActorSystem) |
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(Configuration... testConfigs) |
Modifier and Type | Method and Description |
---|---|
protected static Collection<Object[]> |
TestBaseUtils.toParameterList(List<Configuration> testConfigs) |
Constructor and Description |
---|
AbstractTestBase(Configuration config) |
ForkableFlinkMiniCluster(Configuration userConfiguration) |
ForkableFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
JavaProgramTestBase(Configuration config) |
Modifier and Type | Method and Description |
---|---|
static <T> T |
InstantiationUtil.readObjectFromConfig(Configuration config,
String key,
ClassLoader cl) |
static void |
InstantiationUtil.writeObjectToConfig(Object o,
Configuration config,
String key) |
Modifier and Type | Method and Description |
---|---|
Configuration |
ApplicationMasterBase.createConfiguration(String curDir,
String dynamicPropertiesEncodedString) |
Configuration |
ApplicationClient.flinkConfig() |
Configuration |
FlinkYarnCluster.getFlinkConfiguration() |
Configuration |
FlinkYarnClientBase.getFlinkConfiguration() |
Modifier and Type | Method and Description |
---|---|
static int |
Utils.calculateHeapSize(int memory,
Configuration conf)
See documentation
|
static Map<String,String> |
Utils.getEnvironmentVariables(String envPrefix,
Configuration flinkConfiguration)
Method to extract environment variables from the flinkConfiguration based on the given prefix String.
|
void |
FlinkYarnClientBase.setFlinkConfiguration(Configuration conf) |
Constructor and Description |
---|
ApplicationClient(Configuration flinkConfig,
LeaderRetrievalService leaderRetrievalService) |
FlinkYarnCluster(org.apache.hadoop.yarn.client.api.YarnClient yarnClient,
org.apache.hadoop.yarn.api.records.ApplicationId appId,
org.apache.hadoop.conf.Configuration hadoopConfig,
Configuration flinkConfig,
org.apache.hadoop.fs.Path sessionFilesDir,
boolean detached)
Create a new Flink on YARN cluster.
|
YarnJobManager(Configuration flinkConfiguration,
ExecutorService executorService,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategy restartStrategy,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.