Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion akka-cluster-sharding/src/main/resources/reference.conf
Original file line number Diff line number Diff line change
Expand Up @@ -193,7 +193,7 @@ akka.cluster.sharding {
# is to have more tolerance for membership changes between write and read.
# The tradeoff of increasing this is that updates will be slower.
# It is more important to increase the `read-majority-plus`.
write-majority-plus = 5
write-majority-plus = 3
# State retrieval when ShardCoordinator is started is required to be read
# from a majority of nodes plus this number of additional nodes. Can also
# be set to "all" to require reads from all nodes. The reason for write/read
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -491,7 +491,8 @@ object Replicator {

/**
* `ReadMajority` but with the given number of `additional` nodes added to the majority count. At most
* all nodes.
* all nodes. Exiting nodes are excluded using `ReadMajorityPlus` because those are typically
* about to be removed and will not be able to respond.
*/
final case class ReadMajorityPlus(timeout: FiniteDuration, additional: Int, minCap: Int = DefaultMajorityMinCap)
extends ReadConsistency {
Expand Down Expand Up @@ -535,7 +536,8 @@ object Replicator {

/**
* `WriteMajority` but with the given number of `additional` nodes added to the majority count. At most
* all nodes.
* all nodes. Exiting nodes are excluded using `WriteMajorityPlus` because those are typically
* about to be removed and will not be able to respond.
*/
final case class WriteMajorityPlus(timeout: FiniteDuration, additional: Int, minCap: Int = DefaultMajorityMinCap)
extends WriteConsistency {
Expand Down Expand Up @@ -1460,6 +1462,9 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
// cluster joining nodes, doesn't contain selfAddress
var joiningNodes: immutable.SortedSet[UniqueAddress] = immutable.SortedSet.empty

// cluster exiting nodes, doesn't contain selfAddress
var exitingNodes: immutable.SortedSet[UniqueAddress] = immutable.SortedSet.empty

// up and weaklyUp nodes, doesn't contain joining and not selfAddress
private def allNodes: immutable.SortedSet[UniqueAddress] = nodes.union(weaklyUpNodes)

Expand Down Expand Up @@ -1501,11 +1506,16 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
// messages after loading durable data.
var replyTo: ActorRef = null

private def nodesForReadWrite(): Vector[UniqueAddress] = {
if (settings.preferOldest)
membersByAge.iterator.map(_.uniqueAddress).toVector
else
nodes.toVector
private def nodesForReadWrite(excludeExiting: Boolean): Vector[UniqueAddress] = {
if (excludeExiting && exitingNodes.nonEmpty) {
if (settings.preferOldest) membersByAge.iterator.collect {
case m if !exitingNodes.contains(m.uniqueAddress) => m.uniqueAddress
}.toVector
else nodes.diff(exitingNodes).toVector
} else {
if (settings.preferOldest) membersByAge.iterator.map(_.uniqueAddress).toVector
else nodes.toVector
}
}

override protected[akka] def aroundReceive(rcv: Actor.Receive, msg: Any): Unit = {
Expand Down Expand Up @@ -1661,6 +1671,7 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
case MemberJoined(m) => receiveMemberJoining(m)
case MemberWeaklyUp(m) => receiveMemberWeaklyUp(m)
case MemberUp(m) => receiveMemberUp(m)
case MemberExited(m) => receiveMemberExiting(m)
case MemberRemoved(m, _) => receiveMemberRemoved(m)
case evt: MemberEvent => receiveOtherMemberEvent(evt.member)
case UnreachableMember(m) => receiveUnreachable(m)
Expand All @@ -1683,14 +1694,18 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
}
replyTo ! reply
} else {
val excludeExiting = consistency match {
case _: ReadMajorityPlus | _: ReadAll => true
case _ => false
}
context.actorOf(
ReadAggregator
.props(
key,
consistency,
req,
selfUniqueAddress,
nodesForReadWrite(),
nodesForReadWrite(excludeExiting),
unreachable,
!settings.preferOldest,
localValue,
Expand Down Expand Up @@ -1783,6 +1798,10 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
// of subsequent updates are in sync on the destination nodes.
// The order is also kept when prefer-oldest is enabled.
val shuffle = !(settings.preferOldest || writeDelta.exists(_.requiresCausalDeliveryOfDeltas))
val excludeExiting = writeConsistency match {
case _: WriteMajorityPlus | _: WriteAll => true
case _ => false
}
val writeAggregator =
context.actorOf(
WriteAggregator
Expand All @@ -1793,7 +1812,7 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
writeConsistency,
req,
selfUniqueAddress,
nodesForReadWrite(),
nodesForReadWrite(excludeExiting),
unreachable,
shuffle,
replyTo,
Expand Down Expand Up @@ -1901,6 +1920,10 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
else
replyTo ! DeleteSuccess(key, req)
} else {
val excludeExiting = consistency match {
case _: WriteMajorityPlus | _: WriteAll => true
case _ => false
}
val writeAggregator =
context.actorOf(
WriteAggregator
Expand All @@ -1911,7 +1934,7 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
consistency,
req,
selfUniqueAddress,
nodesForReadWrite(),
nodesForReadWrite(excludeExiting),
unreachable,
!settings.preferOldest,
replyTo,
Expand Down Expand Up @@ -2322,6 +2345,10 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
}
}

def receiveMemberExiting(m: Member): Unit =
if (matchingRole(m) && m.address != selfAddress)
exitingNodes += m.uniqueAddress
Comment thread
johanandren marked this conversation as resolved.

def receiveMemberRemoved(m: Member): Unit = {
if (m.address == selfAddress)
context.stop(self)
Expand All @@ -2332,6 +2359,7 @@ final class Replicator(settings: ReplicatorSettings) extends Actor with ActorLog
nodes -= m.uniqueAddress
weaklyUpNodes -= m.uniqueAddress
joiningNodes -= m.uniqueAddress
exitingNodes -= m.uniqueAddress
removedNodes = removedNodes.updated(m.uniqueAddress, allReachableClockTime)
unreachable -= m.uniqueAddress
if (settings.preferOldest)
Expand Down
14 changes: 8 additions & 6 deletions akka-docs/src/main/paradox/typed/distributed-data.md
Original file line number Diff line number Diff line change
Expand Up @@ -229,9 +229,10 @@ at least **N/2 + 1** replicas, where N is the number of nodes in the cluster
(or cluster role group)
* `WriteMajorityPlus` is like `WriteMajority` but with the given number of `additional` nodes added
to the majority count. At most all nodes. This gives better tolerance for membership changes between
writes and reads.
* `WriteAll` the value will immediately be written to all nodes in the cluster
(or all nodes in the cluster role group)
writes and reads. Exiting nodes are excluded using `WriteMajorityPlus` because those are typically about to be removed
and will not be able to respond.
* `WriteAll` the value will immediately be written to all nodes in the cluster (or all nodes in the cluster role group).
Exiting nodes are excluded using `WriteAll` because those are typically about to be removed and will not be able to respond.

When you specify to write to `n` out of `x` nodes, the update will first replicate to `n` nodes.
If there are not enough Acks after a 1/5th of the timeout, the update will be replicated to `n` other
Expand Down Expand Up @@ -263,9 +264,10 @@ at least **N/2 + 1** replicas, where N is the number of nodes in the cluster
(or cluster role group)
* `ReadMajorityPlus` is like `ReadMajority` but with the given number of `additional` nodes added
to the majority count. At most all nodes. This gives better tolerance for membership changes between
writes and reads.
* `ReadAll` the value will be read and merged from all nodes in the cluster
(or all nodes in the cluster role group)
writes and reads. Exiting nodes are excluded using `ReadMajorityPlus` because those are typically about to be
removed and will not be able to respond.
* `ReadAll` the value will be read and merged from all nodes in the cluster (or all nodes in the cluster role group).
Exiting nodes are excluded using `ReadAll` because those are typically about to be removed and will not be able to respond.

Note that `ReadMajority` and `ReadMajorityPlus` have a `minCap` parameter that is useful to specify to achieve
better safety for small clusters.
Expand Down