public class ForkableFlinkMiniCluster extends LocalFlinkMiniCluster
Constructor and Description |
---|
ForkableFlinkMiniCluster(Configuration userConfiguration) |
ForkableFlinkMiniCluster(Configuration userConfiguration,
boolean singleActorSystem) |
Modifier and Type | Method and Description |
---|---|
Configuration |
generateConfiguration(Configuration userConfiguration) |
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 |
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, initializeIOFormatClasses, setMemory, stopJob
awaitTermination, clearLeader, configuration, createLeaderRetrievalService, executionContext, futureLock, getJobManagerAkkaConfig, getJobManagersAsJava, getLeaderGateway, getLeaderGatewayFuture, getLeaderIndex, getLeaderIndexFuture, getNumberOfJobManagers, getTaskManagerAkkaConfig, getTaskManagers, getTaskManagersAsJava, handleError, hostname, jobManagerActors, jobManagerActorSystems, leaderGateway, leaderIndex, leaderRetrievalService, LOG, notifyLeaderAddress, numJobManagers, numTaskManagers, recoveryMode, running, setDefaultCiConfig, shutdown, shutdownJobClientActorSystem, start, startJobClientActorSystem, startJobManagerActorSystem, 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 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 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.