Skip to content
This repository has been archived by the owner on Apr 18, 2024. It is now read-only.

Rolling update steps for Classic -> Typed Sharding #110

Open
wants to merge 9 commits into
base: 2.6
Choose a base branch
from
Open
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
1 change: 1 addition & 0 deletions akka-sample-sharding-java/build.sbt
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ val `akka-sample-sharding-java` = project
javacOptions in doc in Compile := Seq("-parameters", "-Xdoclint:none"),
libraryDependencies ++= Seq(
"com.typesafe.akka" %% "akka-cluster-sharding" % akkaVersion,
"com.typesafe.akka" %% "akka-cluster-sharding-typed" % akkaVersion,
"com.typesafe.akka" %% "akka-serialization-jackson" % akkaVersion,
"org.scalatest" %% "scalatest" % "3.0.7" % Test
),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,10 +32,12 @@ public String toString() {

public static class GetTemperature implements Command {
public final int deviceId;
public final ActorRef replyTo;

@JsonCreator
public GetTemperature(int deviceId) {
public GetTemperature(int deviceId, ActorRef replyTo) {
this.deviceId = deviceId;
this.replyTo = replyTo;
}
}

Expand Down Expand Up @@ -77,13 +79,18 @@ private void receiveRecordTemperature(RecordTemperature cmd) {
}

private void receiveGetTemperature(GetTemperature cmd) {
Temperature reply;
if (temperatures.isEmpty()) {
getSender().tell(new Temperature(cmd.deviceId, Double.NaN,
Double.NaN, 0), getSelf());
reply = new Temperature(cmd.deviceId, Double.NaN, Double.NaN, 0);
} else {
getSender().tell(new Temperature(cmd.deviceId, average(temperatures),
temperatures.get(temperatures.size() - 1), temperatures.size()), getSelf());
reply = new Temperature(cmd.deviceId, average(temperatures),
temperatures.get(temperatures.size() - 1), temperatures.size());
}

if (cmd.replyTo == null)
getSender().tell(reply, getSelf());
else
cmd.replyTo.tell(reply, getSelf());
}

private double sum(List<Double> values) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import java.time.Duration;
import java.util.Random;
import java.util.concurrent.CompletionStage;

import akka.actor.AbstractActorWithTimers;

Expand All @@ -13,6 +14,8 @@
import akka.cluster.sharding.ClusterSharding;
import akka.cluster.sharding.ClusterShardingSettings;
import akka.cluster.sharding.ShardRegion;
import akka.cluster.sharding.typed.ShardingEnvelope;
import akka.pattern.Patterns;

public class Devices extends AbstractActorWithTimers {

Expand All @@ -31,7 +34,13 @@ public String entityId(Object message) {
return String.valueOf(((Device.RecordTemperature) message).deviceId);
else if (message instanceof Device.GetTemperature)
return String.valueOf(((Device.GetTemperature) message).deviceId);
else
else if (message instanceof ShardingEnvelope) {
ShardingEnvelope envelope = (ShardingEnvelope) message;
if (envelope.message() instanceof Device.RecordTemperature)
return String.valueOf(((Device.RecordTemperature) envelope.message()).deviceId);
else
return null;
} else
return null;
}
};
Expand Down Expand Up @@ -84,8 +93,15 @@ private void receiveUpdateDevice() {
}

private void receiveReadTemperatures() {
for (int id = 0; id < numberOfDevices; id++) {
deviceRegion.tell(new Device.GetTemperature(id), getSelf());
for (int deviceId = 0; deviceId < numberOfDevices; deviceId++) {
if (deviceId >= 40) {
final int id = deviceId;
CompletionStage<Object> reply = Patterns.askWithReplyTo(deviceRegion, replyTo ->
new Device.GetTemperature(id, replyTo), Duration.ofSeconds(3));
Patterns.pipe(reply, getContext().getDispatcher()).to(getSelf());
} else {
deviceRegion.tell(new Device.GetTemperature(deviceId, getSelf()), getSelf());
}
}
}

Expand Down
3 changes: 3 additions & 0 deletions akka-sample-sharding-java/src/main/resources/application.conf
Original file line number Diff line number Diff line change
Expand Up @@ -27,3 +27,6 @@ akka {
}

}

akka.cluster.sharding.number-of-shards = 100

5 changes: 3 additions & 2 deletions akka-sample-sharding-scala/build.sbt
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import com.typesafe.sbt.SbtMultiJvm.multiJvmSettings
import com.typesafe.sbt.SbtMultiJvm.MultiJvmKeys.MultiJvm

val akkaVersion = "2.6.0-M4"
val akkaVersion = "2.6.0"

lazy val `akka-sample-sharding-scala` = project
.in(file("."))
Expand All @@ -20,8 +20,9 @@ lazy val `akka-sample-sharding-scala` = project
javaOptions in run ++= Seq("-Xms128m", "-Xmx1024m"),
libraryDependencies ++= Seq(
"com.typesafe.akka" %% "akka-cluster-sharding" % akkaVersion,
"com.typesafe.akka" %% "akka-cluster-sharding-typed" % akkaVersion,
"com.typesafe.akka" %% "akka-serialization-jackson" % akkaVersion,
"org.scalatest" %% "scalatest" % "3.0.7" % Test
"org.scalatest" %% "scalatest" % "3.0.8" % Test
),
mainClass in (Compile, run) := Some("sample.sharding.ShardingApp"),
// disable parallel tests
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,9 @@ akka {
seed-nodes = [
"akka://[email protected]:2551",
"akka://[email protected]:2552"]

# auto downing is NOT safe for production deployments.
# you may want to use it during development, read more about it in the docs.
auto-down-unreachable-after = 10s
}

}

akka.cluster.sharding.number-of-shards = 100

18 changes: 18 additions & 0 deletions akka-sample-sharding-scala/src/main/resources/logback.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
<configuration>
<appender name="STDOUT" target="System.out" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>[%date{ISO8601}] [%level] [%logger] [%thread] [%X{akkaSource}] - %msg%n</pattern>
</encoder>
</appender>

<appender name="ASYNC" class="ch.qos.logback.classic.AsyncAppender">
<queueSize>1024</queueSize>
<neverBlock>true</neverBlock>
<appender-ref ref="STDOUT" />
</appender>

<root level="INFO">
<appender-ref ref="ASYNC"/>
</root>

</configuration>
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ object Device {
case class RecordTemperature(deviceId: Int, temperature: Double)
extends Command

case class GetTemperature(deviceId: Int) extends Command
case class GetTemperature(deviceId: Int, replyTo: ActorRef) extends Command

case class Temperature(deviceId: Int,
average: Double,
Expand All @@ -38,13 +38,17 @@ class Device extends Actor with ActorLogging {
)
context.become(counting(temperatures))

case GetTemperature(id) =>
case GetTemperature(id, replyTo) =>
val reply =
if (values.isEmpty)
Temperature(id, Double.NaN, Double.NaN, 0)
else
Temperature(id, average(values), values.last, values.size)
sender() ! reply

if (replyTo == null)
sender() ! reply
else
replyTo ! reply

}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,10 @@ import scala.util.Random

import akka.actor._
import akka.cluster.sharding._
import akka.cluster.sharding.typed.ShardingEnvelope
import akka.pattern.extended.ask // note extended.ask
import akka.pattern.pipe
import akka.util.Timeout

object Devices {
// Update a random device
Expand All @@ -21,19 +25,25 @@ class Devices extends Actor with ActorLogging with Timers {

private val extractEntityId: ShardRegion.ExtractEntityId = {
case msg @ Device.RecordTemperature(id, _) => (id.toString, msg)
case msg @ Device.GetTemperature(id) => (id.toString, msg)
case msg @ Device.GetTemperature(id, _) => (id.toString, msg)
case ShardingEnvelope(_, msg @ Device.RecordTemperature(id, _)) =>
(id.toString, msg)
}

private val numberOfShards = 100

private val extractShardId: ShardRegion.ExtractShardId = {
case Device.RecordTemperature(id, _) =>
(math.abs(id.hashCode) % numberOfShards).toString
case Device.GetTemperature(id) =>
case Device.GetTemperature(id, _) =>
(math.abs(id.hashCode) % numberOfShards).toString
// Needed if you want to use 'remember entities':
case ShardRegion.StartEntity(id) =>
(math.abs(id.hashCode) % numberOfShards).toString
case ShardingEnvelope(_, Device.RecordTemperature(id, _)) =>
(math.abs(id.hashCode) % numberOfShards).toString
case ShardingEnvelope(_, Device.GetTemperature(id, _)) =>
(math.abs(id.hashCode) % numberOfShards).toString
}

val deviceRegion: ActorRef = ClusterSharding(context.system).start(
Expand Down Expand Up @@ -64,7 +74,14 @@ class Devices extends Actor with ActorLogging with Timers {

case ReadTemperatures =>
(0 to numberOfDevices).foreach { deviceId =>
deviceRegion ! Device.GetTemperature(deviceId)
if (deviceId >= 40) {
import context.dispatcher
implicit val timeout = Timeout(3.seconds)
deviceRegion
.ask(replyTo => Device.GetTemperature(deviceId, replyTo))
.pipeTo(self)
} else
deviceRegion ! Device.GetTemperature(deviceId, self)
}

case temp: Device.Temperature =>
Expand Down
18 changes: 18 additions & 0 deletions akka-sample-sharding-typed-java/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
*#
*.iml
*.ipr
*.iws
*.pyc
*.tm.epoch
*.vim
*-shim.sbt
.idea/
/project/plugins/project
project/boot
target/
/logs
.cache
.classpath
.project
.settings
native/
121 changes: 121 additions & 0 deletions akka-sample-sharding-typed-java/COPYING
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
Creative Commons Legal Code

CC0 1.0 Universal

CREATIVE COMMONS CORPORATION IS NOT A LAW FIRM AND DOES NOT PROVIDE
LEGAL SERVICES. DISTRIBUTION OF THIS DOCUMENT DOES NOT CREATE AN
ATTORNEY-CLIENT RELATIONSHIP. CREATIVE COMMONS PROVIDES THIS
INFORMATION ON AN "AS-IS" BASIS. CREATIVE COMMONS MAKES NO WARRANTIES
REGARDING THE USE OF THIS DOCUMENT OR THE INFORMATION OR WORKS
PROVIDED HEREUNDER, AND DISCLAIMS LIABILITY FOR DAMAGES RESULTING FROM
THE USE OF THIS DOCUMENT OR THE INFORMATION OR WORKS PROVIDED
HEREUNDER.

Statement of Purpose

The laws of most jurisdictions throughout the world automatically confer
exclusive Copyright and Related Rights (defined below) upon the creator
and subsequent owner(s) (each and all, an "owner") of an original work of
authorship and/or a database (each, a "Work").

Certain owners wish to permanently relinquish those rights to a Work for
the purpose of contributing to a commons of creative, cultural and
scientific works ("Commons") that the public can reliably and without fear
of later claims of infringement build upon, modify, incorporate in other
works, reuse and redistribute as freely as possible in any form whatsoever
and for any purposes, including without limitation commercial purposes.
These owners may contribute to the Commons to promote the ideal of a free
culture and the further production of creative, cultural and scientific
works, or to gain reputation or greater distribution for their Work in
part through the use and efforts of others.

For these and/or other purposes and motivations, and without any
expectation of additional consideration or compensation, the person
associating CC0 with a Work (the "Affirmer"), to the extent that he or she
is an owner of Copyright and Related Rights in the Work, voluntarily
elects to apply CC0 to the Work and publicly distribute the Work under its
terms, with knowledge of his or her Copyright and Related Rights in the
Work and the meaning and intended legal effect of CC0 on those rights.

1. Copyright and Related Rights. A Work made available under CC0 may be
protected by copyright and related or neighboring rights ("Copyright and
Related Rights"). Copyright and Related Rights include, but are not
limited to, the following:

i. the right to reproduce, adapt, distribute, perform, display,
communicate, and translate a Work;
ii. moral rights retained by the original author(s) and/or performer(s);
iii. publicity and privacy rights pertaining to a person's image or
likeness depicted in a Work;
iv. rights protecting against unfair competition in regards to a Work,
subject to the limitations in paragraph 4(a), below;
v. rights protecting the extraction, dissemination, use and reuse of data
in a Work;
vi. database rights (such as those arising under Directive 96/9/EC of the
European Parliament and of the Council of 11 March 1996 on the legal
protection of databases, and under any national implementation
thereof, including any amended or successor version of such
directive); and
vii. other similar, equivalent or corresponding rights throughout the
world based on applicable law or treaty, and any national
implementations thereof.

2. Waiver. To the greatest extent permitted by, but not in contravention
of, applicable law, Affirmer hereby overtly, fully, permanently,
irrevocably and unconditionally waives, abandons, and surrenders all of
Affirmer's Copyright and Related Rights and associated claims and causes
of action, whether now known or unknown (including existing as well as
future claims and causes of action), in the Work (i) in all territories
worldwide, (ii) for the maximum duration provided by applicable law or
treaty (including future time extensions), (iii) in any current or future
medium and for any number of copies, and (iv) for any purpose whatsoever,
including without limitation commercial, advertising or promotional
purposes (the "Waiver"). Affirmer makes the Waiver for the benefit of each
member of the public at large and to the detriment of Affirmer's heirs and
successors, fully intending that such Waiver shall not be subject to
revocation, rescission, cancellation, termination, or any other legal or
equitable action to disrupt the quiet enjoyment of the Work by the public
as contemplated by Affirmer's express Statement of Purpose.

3. Public License Fallback. Should any part of the Waiver for any reason
be judged legally invalid or ineffective under applicable law, then the
Waiver shall be preserved to the maximum extent permitted taking into
account Affirmer's express Statement of Purpose. In addition, to the
extent the Waiver is so judged Affirmer hereby grants to each affected
person a royalty-free, non transferable, non sublicensable, non exclusive,
irrevocable and unconditional license to exercise Affirmer's Copyright and
Related Rights in the Work (i) in all territories worldwide, (ii) for the
maximum duration provided by applicable law or treaty (including future
time extensions), (iii) in any current or future medium and for any number
of copies, and (iv) for any purpose whatsoever, including without
limitation commercial, advertising or promotional purposes (the
"License"). The License shall be deemed effective as of the date CC0 was
applied by Affirmer to the Work. Should any part of the License for any
reason be judged legally invalid or ineffective under applicable law, such
partial invalidity or ineffectiveness shall not invalidate the remainder
of the License, and in such case Affirmer hereby affirms that he or she
will not (i) exercise any of his or her remaining Copyright and Related
Rights in the Work or (ii) assert any associated claims and causes of
action with respect to the Work, in either case contrary to Affirmer's
express Statement of Purpose.

4. Limitations and Disclaimers.

a. No trademark or patent rights held by Affirmer are waived, abandoned,
surrendered, licensed or otherwise affected by this document.
b. Affirmer offers the Work as-is and makes no representations or
warranties of any kind concerning the Work, express, implied,
statutory or otherwise, including without limitation warranties of
title, merchantability, fitness for a particular purpose, non
infringement, or the absence of latent or other defects, accuracy, or
the present or absence of errors, whether or not discoverable, all to
the greatest extent permissible under applicable law.
c. Affirmer disclaims responsibility for clearing rights of other persons
that may apply to the Work or any use thereof, including without
limitation any person's Copyright and Related Rights in the Work.
Further, Affirmer disclaims responsibility for obtaining any necessary
consents, permissions or other rights required for any use of the
Work.
d. Affirmer understands and acknowledges that Creative Commons is not a
party to this document and has no duty or obligation with respect to
this CC0 or use of the Work.
10 changes: 10 additions & 0 deletions akka-sample-sharding-typed-java/LICENSE
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
Akka sample by Lightbend

Licensed under Public Domain (CC0)

To the extent possible under law, the person who associated CC0 with
this Template has waived all copyright and related or neighboring
rights to this Template.

You should have received a copy of the CC0 legalcode along with this
work. If not, see <http://creativecommons.org/publicdomain/zero/1.0/>.
Loading