public class ForkableFlinkMiniCluster extends LocalFlinkMiniCluster
Constructor and Description |
---|
ForkableFlinkMiniCluster(Configuration userConfiguration) |
ForkableFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
Modifier and Type | Method and Description |
---|---|
static String |
DEFAULT_MINICLUSTER_AKKA_ASK_TIMEOUT() |
Configuration |
generateConfiguration(Configuration userConfiguration) |
static scala.concurrent.duration.FiniteDuration |
MAX_RESTART_DURATION() |
void |
restartLeadingJobManager() |
void |
restartTaskManager(int index) |
void |
start() |
static ForkableFlinkMiniCluster |
startCluster(int numSlots,
int numTaskManagers,
String timeout) |
akka.actor.ActorRef |
startJobManager(int index,
akka.actor.ActorSystem actorSystem) |
akka.actor.ActorRef |
startResourceManager(int index,
akka.actor.ActorSystem system) |
akka.actor.ActorRef |
startTaskManager(int index,
akka.actor.ActorSystem system) |
void |
stop() |
void |
waitForTaskManagersToBeRegisteredAtJobManager(akka.actor.ActorRef jobManager) |
scala.Option<org.apache.curator.test.TestingCluster> |
zookeeperCluster() |
currentlyRunningJobs, getArchiveName, getCurrentlyRunningJobsJava, getDefaultConfig, getJobManagerName, getLeaderRPCPort, getResourceManagerName, initializeIOFormatClasses, setMemory, stopJob
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, startJobClientActorSystem, startJobManagerActorSystem, startResourceManagerActorSystem, startTaskManagerActorSystem, startWebServer, submitJobAndWait, submitJobAndWait, submitJobAndWait, submitJobDetached, taskManagerActors, taskManagerActorSystems, timeout, userConfiguration, useSingleActorSystem, waitForTaskManagersToBeRegistered, waitForTaskManagersToBeRegistered, webMonitor
public ForkableFlinkMiniCluster(Configuration userConfiguration, boolean singleActorSystem)
public ForkableFlinkMiniCluster(Configuration userConfiguration)
public static scala.concurrent.duration.FiniteDuration MAX_RESTART_DURATION()
public static String DEFAULT_MINICLUSTER_AKKA_ASK_TIMEOUT()
public static ForkableFlinkMiniCluster startCluster(int numSlots, int numTaskManagers, String timeout)
public scala.Option<org.apache.curator.test.TestingCluster> zookeeperCluster()
public Configuration generateConfiguration(Configuration userConfiguration)
generateConfiguration
in class LocalFlinkMiniCluster
public akka.actor.ActorRef startJobManager(int index, akka.actor.ActorSystem actorSystem)
startJobManager
in class LocalFlinkMiniCluster
public akka.actor.ActorRef startResourceManager(int index, akka.actor.ActorSystem system)
startResourceManager
in class LocalFlinkMiniCluster
public akka.actor.ActorRef startTaskManager(int index, akka.actor.ActorSystem system)
startTaskManager
in class LocalFlinkMiniCluster
public void restartLeadingJobManager()
public void restartTaskManager(int index)
public void start()
start
in class FlinkMiniCluster
public void stop()
stop
in class FlinkMiniCluster
public void waitForTaskManagersToBeRegisteredAtJobManager(akka.actor.ActorRef jobManager)
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.