public class LocalFlinkMiniCluster extends FlinkMiniCluster
TaskManager
s and the JobManager
in the same
JVM. It extends the FlinkMiniCluster
by having convenience functions to setup Flink's
configuration and implementations to create JobManager
and TaskManager
.
Constructor and Description |
---|
LocalFlinkMiniCluster(Configuration userConfiguration) |
LocalFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
Modifier and Type | Method and Description |
---|---|
scala.collection.Iterable<JobID> |
currentlyRunningJobs() |
Configuration |
generateConfiguration(Configuration userConfiguration) |
protected String |
getArchiveName(int index) |
List<JobID> |
getCurrentlyRunningJobsJava() |
Configuration |
getDefaultConfig() |
protected String |
getJobManagerName(int index) |
int |
getLeaderRPCPort() |
protected String |
getResourceManagerName(int index) |
void |
initializeIOFormatClasses(Configuration configuration) |
void |
setMemory(Configuration config) |
akka.actor.ActorRef |
startJobManager(int index,
akka.actor.ActorSystem system) |
akka.actor.ActorRef |
startResourceManager(int index,
akka.actor.ActorSystem system) |
akka.actor.ActorRef |
startTaskManager(int index,
akka.actor.ActorSystem system) |
void |
stopJob(JobID id) |
awaitTermination, clearLeader, configuration, createLeaderRetrievalService, executionContext, futureExecutor, futureLock, getJobManagerAkkaConfig, getJobManagersAsJava, getLeaderGateway, getLeaderGatewayFuture, getLeaderIndex, getLeaderIndexFuture, getNumberOfJobManagers, getNumberOfResourceManagers, getResourceManagerAkkaConfig, getTaskManagerAkkaConfig, getTaskManagers, getTaskManagersAsJava, handleError, hostname, ioExecutor, jobManagerActors, jobManagerActorSystems, jobManagerLeaderRetrievalService, leaderGateway, leaderIndex, LOG, notifyLeaderAddress, numJobManagers, numTaskManagers, recoveryMode, resourceManagerActors, resourceManagerActorSystems, running, setDefaultCiConfig, shutdown, shutdownJobClientActorSystem, start, start, startJobClientActorSystem, startJobManagerActorSystem, startResourceManagerActorSystem, startTaskManagerActorSystem, startWebServer, stop, submitJobAndWait, submitJobAndWait, submitJobAndWait, submitJobDetached, taskManagerActors, taskManagerActorSystems, timeout, userConfiguration, useSingleActorSystem, waitForTaskManagersToBeRegistered, waitForTaskManagersToBeRegistered, webMonitor
public LocalFlinkMiniCluster(Configuration userConfiguration, boolean singleActorSystem)
public LocalFlinkMiniCluster(Configuration userConfiguration)
public Configuration generateConfiguration(Configuration userConfiguration)
generateConfiguration
in class FlinkMiniCluster
public akka.actor.ActorRef startJobManager(int index, akka.actor.ActorSystem system)
startJobManager
in class FlinkMiniCluster
public akka.actor.ActorRef startResourceManager(int index, akka.actor.ActorSystem system)
startResourceManager
in class FlinkMiniCluster
public akka.actor.ActorRef startTaskManager(int index, akka.actor.ActorSystem system)
startTaskManager
in class FlinkMiniCluster
public int getLeaderRPCPort()
public void initializeIOFormatClasses(Configuration configuration)
public void setMemory(Configuration config)
public Configuration getDefaultConfig()
protected String getJobManagerName(int index)
protected String getResourceManagerName(int index)
protected String getArchiveName(int index)
public scala.collection.Iterable<JobID> currentlyRunningJobs()
public void stopJob(JobID id)
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.