Class ShardCoordinator
- java.lang.Object
-
- org.apache.pekko.cluster.sharding.ShardCoordinator
-
- Direct Known Subclasses:
PersistentShardCoordinator
public abstract class ShardCoordinator extends java.lang.Object implements Actor, Timers
Singleton coordinator that decides where to allocate shards.- See Also:
ClusterSharding extension
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description static class
ShardCoordinator.AbstractShardAllocationStrategy
Java API: Java implementations of custom shard allocation and rebalancing logic used by theShardCoordinator
should extend this abstract class and implement the two methods.static interface
ShardCoordinator.ActorSystemDependentAllocationStrategy
Shard allocation strategy where start is called by the shard coordinator before any calls to rebalance or allocate shard.static class
ShardCoordinator.Internal$
INTERNAL APIstatic class
ShardCoordinator.LeastShardAllocationStrategy
Usepekko.cluster.sharding.ShardCoordinator.ShardAllocationStrategy.leastShardAllocationStrategy
instead.static class
ShardCoordinator.RebalanceWorker$
static interface
ShardCoordinator.ShardAllocationStrategy
Interface of the pluggable shard allocation and rebalancing logic used by theShardCoordinator
.static class
ShardCoordinator.ShardAllocationStrategy$
static interface
ShardCoordinator.StartableAllocationStrategy
Shard allocation strategy where start is called by the shard coordinator before any calls to rebalance or allocate shard.-
Nested classes/interfaces inherited from interface org.apache.pekko.actor.Actor
Actor.emptyBehavior$, Actor.ignoringBehavior$
-
-
Constructor Summary
Constructors Constructor Description ShardCoordinator(ClusterShardingSettings settings, ShardCoordinator.ShardAllocationStrategy allocationStrategy)
-
Method Summary
All Methods Static Methods Instance Methods Abstract Methods Concrete Methods Modifier and Type Method Description scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit>
active()
scala.collection.immutable.Set<ActorRef>
aliveRegions()
void
aliveRegions_$eq(scala.collection.immutable.Set<ActorRef> x$1)
void
allocateShardHomesForRememberEntities()
boolean
allRegionsRegistered()
void
allRegionsRegistered_$eq(boolean x$1)
Cluster
cluster()
ActorContext
context()
Scala API: Stores the context for this actor, including self, and sender.void
continueGetShardHome(java.lang.String shard, ActorRef region, ActorRef getShardHomeSender)
void
continueRebalance(scala.collection.immutable.Set<java.lang.String> shards)
scala.collection.immutable.Set<ActorRef>
gracefulShutdownInProgress()
void
gracefulShutdownInProgress_$eq(scala.collection.immutable.Set<ActorRef> x$1)
boolean
handleGetShardHome(java.lang.String shard)
boolean
hasAllRegionsRegistered()
boolean
isMember(ActorRef region)
static ShardCoordinator.ShardAllocationStrategy
leastShardAllocationStrategy(int absoluteLimit, double relativeLimit)
Java API:ShardAllocationStrategy
that allocates new shards to theShardRegion
(node) with least number of previously allocated shards.MarkerLoggingAdapter
log()
int
minMembers()
protected void
org$apache$pekko$actor$Actor$_setter_$context_$eq(ActorContext x$1)
Scala API: Stores the context for this actor, including self, and sender.protected void
org$apache$pekko$actor$Actor$_setter_$self_$eq(ActorRef x$1)
The 'self' field holds the ActorRef for this actor.void
postStop()
User overridable callback.boolean
preparingForShutdown()
void
preparingForShutdown_$eq(boolean x$1)
void
preStart()
User overridable callback.scala.collection.immutable.Map<java.lang.String,scala.collection.immutable.Set<ActorRef>>
rebalanceInProgress()
void
rebalanceInProgress_$eq(scala.collection.immutable.Map<java.lang.String,scala.collection.immutable.Set<ActorRef>> x$1)
scala.collection.immutable.Set<ActorRef>
rebalanceWorkers()
void
rebalanceWorkers_$eq(scala.collection.immutable.Set<ActorRef> x$1)
scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit>
receiveTerminated()
void
regionProxyTerminated(ActorRef ref)
void
regionTerminated(ActorRef ref)
scala.collection.immutable.Set<ActorRef>
regionTerminationInProgress()
void
regionTerminationInProgress_$eq(scala.collection.immutable.Set<ActorRef> x$1)
scala.concurrent.duration.FiniteDuration
removalMargin()
ActorRef
self()
The 'self' field holds the ActorRef for this actor.void
sendHostShardMsg(java.lang.String shard, ActorRef region)
void
shutdownShards(ActorRef shuttingDownRegion, scala.collection.immutable.Set<java.lang.String> shards)
scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit>
shuttingDown()
org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.State
state()
void
state_$eq(org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.State x$1)
void
stateInitialized()
protected abstract java.lang.String
typeName()
scala.collection.immutable.Map<java.lang.String,Cancellable>
unAckedHostShards()
void
unAckedHostShards_$eq(scala.collection.immutable.Map<java.lang.String,Cancellable> x$1)
protected abstract void
unstashOneGetShardHomeRequest()
abstract <E extends org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.DomainEvent>
voidupdate(E evt, scala.Function1<E,scala.runtime.BoxedUnit> f)
boolean
waitingForLocalRegionToTerminate()
void
waitingForLocalRegionToTerminate_$eq(boolean x$1)
void
watchStateActors()
-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
-
Methods inherited from interface org.apache.pekko.actor.Actor
aroundPostRestart, aroundPreStart, postRestart, preRestart, receive, sender, supervisorStrategy, unhandled
-
Methods inherited from interface org.apache.pekko.actor.Timers
actorCell, aroundPostStop, aroundPreRestart, aroundReceive, super$aroundPostStop, super$aroundPreRestart, super$aroundReceive, timers
-
-
-
-
Constructor Detail
-
ShardCoordinator
public ShardCoordinator(ClusterShardingSettings settings, ShardCoordinator.ShardAllocationStrategy allocationStrategy)
-
-
Method Detail
-
leastShardAllocationStrategy
public static ShardCoordinator.ShardAllocationStrategy leastShardAllocationStrategy(int absoluteLimit, double relativeLimit)
Java API:ShardAllocationStrategy
that allocates new shards to theShardRegion
(node) with least number of previously allocated shards.When a node is added to the cluster the shards on the existing nodes will be rebalanced to the new node. The
LeastShardAllocationStrategy
picks shards for rebalancing from theShardRegion
s with most number of previously allocated shards. They will then be allocated to theShardRegion
with least number of previously allocated shards, i.e. new members in the cluster. The amount of shards to rebalance in each round can be limited to make it progress slower since rebalancing too many shards at the same time could result in additional load on the system. For example, causing many Event Sourced entites to be started at the same time.It will not rebalance when there is already an ongoing rebalance in progress.
- Parameters:
absoluteLimit
- the maximum number of shards that will be rebalanced in one rebalance roundrelativeLimit
- fraction (< 1.0) of total number of (known) shards that will be rebalanced in one rebalance round
-
context
public ActorContext context()
Description copied from interface:Actor
Scala API: Stores the context for this actor, including self, and sender. It is implicit to support operations such asforward
.WARNING: Only valid within the Actor itself, so do not close over it and publish it to other threads!
pekko.actor.ActorContext
is the Scala API.getContext
returns apekko.actor.AbstractActor.ActorContext
, which is the Java API of the actor context.
-
self
public final ActorRef self()
Description copied from interface:Actor
The 'self' field holds the ActorRef for this actor. Can be used to send messages to itself:self ! message
-
org$apache$pekko$actor$Actor$_setter_$context_$eq
protected void org$apache$pekko$actor$Actor$_setter_$context_$eq(ActorContext x$1)
Description copied from interface:Actor
Scala API: Stores the context for this actor, including self, and sender. It is implicit to support operations such asforward
.WARNING: Only valid within the Actor itself, so do not close over it and publish it to other threads!
pekko.actor.ActorContext
is the Scala API.getContext
returns apekko.actor.AbstractActor.ActorContext
, which is the Java API of the actor context.- Specified by:
org$apache$pekko$actor$Actor$_setter_$context_$eq
in interfaceActor
-
org$apache$pekko$actor$Actor$_setter_$self_$eq
protected final void org$apache$pekko$actor$Actor$_setter_$self_$eq(ActorRef x$1)
Description copied from interface:Actor
The 'self' field holds the ActorRef for this actor. Can be used to send messages to itself:self ! message
- Specified by:
org$apache$pekko$actor$Actor$_setter_$self_$eq
in interfaceActor
-
log
public MarkerLoggingAdapter log()
-
cluster
public Cluster cluster()
-
removalMargin
public scala.concurrent.duration.FiniteDuration removalMargin()
-
minMembers
public int minMembers()
-
allRegionsRegistered
public boolean allRegionsRegistered()
-
allRegionsRegistered_$eq
public void allRegionsRegistered_$eq(boolean x$1)
-
state
public org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.State state()
-
state_$eq
public void state_$eq(org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.State x$1)
-
preparingForShutdown
public boolean preparingForShutdown()
-
preparingForShutdown_$eq
public void preparingForShutdown_$eq(boolean x$1)
-
rebalanceInProgress
public scala.collection.immutable.Map<java.lang.String,scala.collection.immutable.Set<ActorRef>> rebalanceInProgress()
-
rebalanceInProgress_$eq
public void rebalanceInProgress_$eq(scala.collection.immutable.Map<java.lang.String,scala.collection.immutable.Set<ActorRef>> x$1)
-
rebalanceWorkers
public scala.collection.immutable.Set<ActorRef> rebalanceWorkers()
-
rebalanceWorkers_$eq
public void rebalanceWorkers_$eq(scala.collection.immutable.Set<ActorRef> x$1)
-
unAckedHostShards
public scala.collection.immutable.Map<java.lang.String,Cancellable> unAckedHostShards()
-
unAckedHostShards_$eq
public void unAckedHostShards_$eq(scala.collection.immutable.Map<java.lang.String,Cancellable> x$1)
-
gracefulShutdownInProgress
public scala.collection.immutable.Set<ActorRef> gracefulShutdownInProgress()
-
gracefulShutdownInProgress_$eq
public void gracefulShutdownInProgress_$eq(scala.collection.immutable.Set<ActorRef> x$1)
-
waitingForLocalRegionToTerminate
public boolean waitingForLocalRegionToTerminate()
-
waitingForLocalRegionToTerminate_$eq
public void waitingForLocalRegionToTerminate_$eq(boolean x$1)
-
aliveRegions
public scala.collection.immutable.Set<ActorRef> aliveRegions()
-
aliveRegions_$eq
public void aliveRegions_$eq(scala.collection.immutable.Set<ActorRef> x$1)
-
regionTerminationInProgress
public scala.collection.immutable.Set<ActorRef> regionTerminationInProgress()
-
regionTerminationInProgress_$eq
public void regionTerminationInProgress_$eq(scala.collection.immutable.Set<ActorRef> x$1)
-
typeName
protected abstract java.lang.String typeName()
-
preStart
public void preStart()
Description copied from interface:Actor
User overridable callback. Is called when an Actor is started. Actors are automatically started asynchronously when created. Empty default implementation.
-
postStop
public void postStop()
Description copied from interface:Actor
User overridable callback. Is called asynchronously after 'actor.stop()' is invoked. Empty default implementation.
-
isMember
public boolean isMember(ActorRef region)
-
active
public scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit> active()
-
handleGetShardHome
public boolean handleGetShardHome(java.lang.String shard)
- Returns:
true
if the message could be handled without state update, i.e. the shard location was known or the request was deferred or ignored
-
receiveTerminated
public scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit> receiveTerminated()
-
update
public abstract <E extends org.apache.pekko.cluster.sharding.ShardCoordinator.Internal.DomainEvent> void update(E evt, scala.Function1<E,scala.runtime.BoxedUnit> f)
-
watchStateActors
public void watchStateActors()
-
stateInitialized
public void stateInitialized()
-
hasAllRegionsRegistered
public boolean hasAllRegionsRegistered()
-
regionTerminated
public void regionTerminated(ActorRef ref)
-
regionProxyTerminated
public void regionProxyTerminated(ActorRef ref)
-
shuttingDown
public scala.PartialFunction<java.lang.Object,scala.runtime.BoxedUnit> shuttingDown()
-
sendHostShardMsg
public void sendHostShardMsg(java.lang.String shard, ActorRef region)
-
allocateShardHomesForRememberEntities
public void allocateShardHomesForRememberEntities()
-
continueGetShardHome
public void continueGetShardHome(java.lang.String shard, ActorRef region, ActorRef getShardHomeSender)
-
unstashOneGetShardHomeRequest
protected abstract void unstashOneGetShardHomeRequest()
-
continueRebalance
public void continueRebalance(scala.collection.immutable.Set<java.lang.String> shards)
-
shutdownShards
public void shutdownShards(ActorRef shuttingDownRegion, scala.collection.immutable.Set<java.lang.String> shards)
-
-