This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-samples.git
The following commit(s) were added to refs/heads/main by this push:
new 56dbe52 use testcontainers to run cassandra (#142)
56dbe52 is described below
commit 56dbe5284e276dc6e09162210d906ded868e4806
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Oct 7 14:32:32 2026 +0100
use testcontainers to run cassandra (#142)
* use testcontainers to run cassandra
* Bind Cassandra test container to port 9042
Motivation:
Testcontainers maps container port 9042 to a random host port, so the
sample nodes (which use the default driver contact point
127.0.0.1:9042) could not reach the Cassandra started by the sample.
Modification:
Bind host port 9042 to the container's 9042, upgrade testcontainers to
2.0.5 (Docker Engine 29 API compatibility), and drop unused imports.
Result:
The samples start Cassandra via Testcontainers and persist events.
Tests:
- distributed-workers-scala: compiled; `cassandra` mode + separate
back-end/front-end/worker nodes processed work, events in pekko.messages
- persistence-dc-scala and persistence-dc-java: compiled; single-JVM mode
replicated thumbs-up across eu-west and eu-central via HTTP
- scalafmt --test passes on changed Scala files
References:
Refs #142
* Fix distributed-workers seed node role and document Docker requirement
Motivation:
The README starts the first seed node on port 7345, but Main treated
ports outside 2000-2999 as worker nodes, so no WorkManager singleton
ran. The port binding used java.util.List.of, which is not available
on Java 8. Creating the journal keyspace in a freshly started Cassandra
container can exceed the driver's default 2s request timeout. The
READMEs did not mention that Docker is required.
Modification:
Treat seed ports 7345 and 7355 as back-end nodes, use
Collections.singletonList, raise the driver request timeout to 10s,
and update the READMEs for Docker and the container-based journal.
Result:
The documented multi-process flow processes work, and the samples
compile and run on Java 8.
Tests:
- All three samples clean-compiled with JDK 8 (Temurin 1.8.0_502)
- distributed-workers-scala on JDK 8: README flow (cassandra, 7345,
3001, 5001 3) joined with correct roles and completed work
- scalafmt --test passes on changed Scala files
References:
Refs #142
* Move distributed-workers back-end port range to 7000-7999
Motivation:
Review feedback asked whether backEndPortRange should stay. It should:
it is the documented way to start additional back-end nodes. But the
documented back-end ports 7345 and 7355 fell outside the 2000-2999
range, which needed a separate set of ports (misnamed backEndSeedPorts,
as only 7345 is a seed node) to start them as back-end nodes.
Modification:
Change the back-end range to 7000-7999, which contains 7345 and 7355,
drop the separate set of ports and fix the comment. Update the README
accordingly.
Result:
A single range decides the back-end role, covering the seed node 7345.
Tests:
- pekko-sample-distributed-workers-scala: sbt compile - pass
References:
Refs #142
* Fork run in distributed-workers so the nodes stay up
Motivation:
With sbt's non-forked run, `sbt run` and the README's `sbt "runMain
worker.Main <port>"` commands finish as soon as main returns, and sbt then
shuts down all the cluster nodes before they can form a cluster or
process any work.
Modification:
Set `run / fork := true`, so sbt waits for the forked JVM, which keeps
running until Ctrl-C.
Result:
`sbt run` keeps the nodes running and processing work until Ctrl-C, which
stops the nodes and the Cassandra container.
Tests:
- sbt run: 6 nodes up, work completed, still running after 95s
- sbt run in a pseudo-terminal, Ctrl-C after 13 jobs done: sbt, the
forked JVM and the Cassandra container all stopped
References:
Refs #142
---
pekko-sample-distributed-workers-scala/README.md | 20 ++++++++------------
pekko-sample-distributed-workers-scala/build.sbt | 6 +++++-
.../src/main/resources/application.conf | 3 +++
.../src/main/scala/worker/Main.scala | 21 +++++++++------------
pekko-sample-persistence-dc-java/README.md | 2 ++
pekko-sample-persistence-dc-java/pom.xml | 6 +++---
.../main/java/sample/persistence/res/MainApp.java | 18 ++++++++++--------
.../src/main/resources/application.conf | 2 ++
pekko-sample-persistence-dc-scala/README.md | 2 ++
pekko-sample-persistence-dc-scala/build.sbt | 7 +------
.../src/main/resources/application.conf | 2 ++
.../main/scala/sample/persistence/res/MainApp.scala | 19 +++++++++++--------
12 files changed, 58 insertions(+), 50 deletions(-)
diff --git a/pekko-sample-distributed-workers-scala/README.md
b/pekko-sample-distributed-workers-scala/README.md
index 1c5c7cc..a7c9d80 100644
--- a/pekko-sample-distributed-workers-scala/README.md
+++ b/pekko-sample-distributed-workers-scala/README.md
@@ -191,7 +191,7 @@ Now that we have covered all the details, we can experiment
with different sets
## Experimenting
-When running the application without parameters it runs a six node cluster
within the same JVM and starts a Apache Cassandra database. It can be more
interesting to run them in separate processes. Open four terminal windows.
+When running the application without parameters it runs a six node cluster
within the same JVM and starts a Apache Cassandra database. The Apache
Cassandra database is started in a Docker container using
[Testcontainers](https://testcontainers.com/), so Docker must be installed and
running. It can be more interesting to run them in separate processes. Open
four terminal windows.
In the first terminal window, start the Apache Cassandra database with the
following command:
@@ -209,7 +209,7 @@ With the database running, go to the second terminal window
and start the first
sbt "runMain worker.Main 7345"
```
-7345 corresponds to the port of the first seed-nodes element in the
configuration. In the log output you see that the cluster node has been started
and changed status to 'Up'.
+7345 corresponds to the port of the first seed-nodes element in the
configuration. Ports 7345 and 7355 start back-end nodes. In the log output you
see that the cluster node has been started and changed status to 'Up'.
In the third terminal window, start the front-end node with the following
command:
@@ -241,35 +241,31 @@ sbt "runMain worker.Main 5001 3"
You can also start more such worker nodes in new terminal windows.
-You can start more cluster back-end nodes using port numbers between 2000-2999.
+You can start more cluster back-end nodes using port numbers between 7000-7999.
```bash
sbt "runMain worker.Main 7355"
```
-The nodes with port 7345 to 2554 are configured to be used as "seed nodes" in
this sample, if you shutdown all or start none of these the other nodes will
not know how to join the cluster. If all four are shut down and 7345 is started
it will join itself and form a new cluster.
+The nodes with ports 7345 and 3000 are configured to be used as "seed nodes"
in this sample, if you shutdown all or start none of these the other nodes will
not know how to join the cluster. If both are shut down and 7345 is started it
will join itself and form a new cluster.
-As long as one of the four nodes is alive the cluster will keep working. You
can read more about this in the [Pekkodocumentation section on seed
nodes](https://pekko.apache.org/docs/pekko/current/scala/cluster-usage.html).
+As long as one of the seed nodes is alive the cluster will keep working. You
can read more about this in the [Pekkodocumentation section on seed
nodes](https://pekko.apache.org/docs/pekko/current/scala/cluster-usage.html).
You can start more cluster front-end nodes using port numbers between
3000-3999:
```bash
-sbt "runMain worker.Main 3002
+sbt "runMain worker.Main 3002"
```
Any port outside these ranges creates a worker node, for which you can also
play around with the number of worker actors on using the second parameter.
```bash
-sbt "runMain worker.Main 5009 4
+sbt "runMain worker.Main 5009 4"
```
## The journal
-The files of the Apache Cassandra database are saved in the target directory
and when you restart the application the state is recovered. You can clean the
state with:
-
-```bash
-sbt clean
-```
+The Apache Cassandra database runs in a Docker container that is removed when
the process that started it stops. While the database is running, restarted
back-end nodes recover their state from the journal. Stopping the database
discards all of the stored state, so the next run starts with an empty journal.
## Next steps
diff --git a/pekko-sample-distributed-workers-scala/build.sbt
b/pekko-sample-distributed-workers-scala/build.sbt
index d5ef2c6..26a3b3d 100644
--- a/pekko-sample-distributed-workers-scala/build.sbt
+++ b/pekko-sample-distributed-workers-scala/build.sbt
@@ -10,6 +10,10 @@ val logbackVersion = "1.3.15"
Global / cancelable := false
+// run in a forked JVM so that sbt keeps the nodes running until Ctrl-C,
instead of
+// stopping them as soon as the main method returns
+run / fork := true
+
libraryDependencies ++= Seq(
"org.apache.pekko" %% "pekko-cluster-typed" % pekkoVersion,
"org.apache.pekko" %% "pekko-persistence-typed" % pekkoVersion,
@@ -17,7 +21,7 @@ libraryDependencies ++= Seq(
"org.apache.pekko" %% "pekko-serialization-jackson" % pekkoVersion,
"org.apache.pekko" %% "pekko-persistence-cassandra" % cassandraPluginVersion,
// this allows us to start cassandra from the sample
- "org.apache.pekko" %% "pekko-persistence-cassandra-launcher" %
cassandraPluginVersion,
+ "org.testcontainers" % "testcontainers-cassandra" % "2.0.5",
"ch.qos.logback" % "logback-classic" % logbackVersion,
// test dependencies
"org.apache.pekko" %% "pekko-actor-testkit-typed" % pekkoVersion % Test,
diff --git
a/pekko-sample-distributed-workers-scala/src/main/resources/application.conf
b/pekko-sample-distributed-workers-scala/src/main/resources/application.conf
index c1ce11f..01a4b2f 100644
--- a/pekko-sample-distributed-workers-scala/src/main/resources/application.conf
+++ b/pekko-sample-distributed-workers-scala/src/main/resources/application.conf
@@ -47,6 +47,9 @@ pekko.persistence.cassandra {
}
}
+# creating the keyspace and tables in a freshly started Cassandra container
can exceed the default 2s timeout
+datastax-java-driver.basic.request.timeout = 10s
+
# Configuration related to the app is in its own namespace
distributed-workers {
# If a workload hasn't finished in this long it
diff --git
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
index 03002fd..898403e 100644
--- a/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
+++ b/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
@@ -1,19 +1,16 @@
package worker
-import java.io.File
import java.util.concurrent.CountDownLatch
import org.apache.pekko.actor.typed.ActorSystem
import org.apache.pekko.actor.typed.eventstream.EventStream
import org.apache.pekko.actor.typed.scaladsl.{ ActorContext, Behaviors }
import org.apache.pekko.cluster.typed.{ Cluster, SelfUp, Subscribe }
-import org.apache.pekko.persistence.cassandra.testkit.CassandraLauncher
import com.typesafe.config.{ Config, ConfigFactory }
object Main {
- // note that 7345 and 7355 are expected to be seed nodes though, even if
- // the back-end starts at 2000
- val backEndPortRange = 2000 to 2999
+ // includes 7345, which is also a seed node (see pekko.cluster.seed-nodes in
application.conf)
+ val backEndPortRange = 7000 to 7999
val frontEndPortRange = 3000 to 3999
@@ -93,16 +90,16 @@ object Main {
* in a real application a pre-existing Apache Cassandra cluster should be
used.
*/
def startCassandraDatabase(): Unit = {
- val databaseDirectory = new File("target/cassandra-db")
- CassandraLauncher.start(
- databaseDirectory,
- CassandraLauncher.DefaultTestConfigResource,
- clean = false,
- port = 9042)
+ import org.testcontainers.cassandra.CassandraContainer
+ import org.testcontainers.utility.DockerImageName
+ val container = new
CassandraContainer(DockerImageName.parse("cassandra:5.0.5"))
+ // bind to the fixed port 9042 so that the sample nodes can connect with
the default driver settings
+ container.setPortBindings(java.util.Collections.singletonList("9042:9042"))
+ container.start()
// shut the cassandra instance down when the JVM stops
sys.addShutdownHook {
- CassandraLauncher.stop()
+ container.stop()
}
}
diff --git a/pekko-sample-persistence-dc-java/README.md
b/pekko-sample-persistence-dc-java/README.md
index 369ba73..7424678 100644
--- a/pekko-sample-persistence-dc-java/README.md
+++ b/pekko-sample-persistence-dc-java/README.md
@@ -3,6 +3,8 @@ pekko-sample-persistence-dc-java
## How to run
+Starting Apache Cassandra with the `cassandra` argument runs it in a Docker
container using [Testcontainers](https://testcontainers.com/), so Docker must
be installed and running.
+
1. Setup Apache Cassandra
* Either start a local Cassandra listening on port 9042 or in terminal 1:
`mvn exec:java -Dexec.mainClass=sample.persistence.res.MainApp
-Dexec.args="cassandra"`
diff --git a/pekko-sample-persistence-dc-java/pom.xml
b/pekko-sample-persistence-dc-java/pom.xml
index f4647f6..54f36ff 100644
--- a/pekko-sample-persistence-dc-java/pom.xml
+++ b/pekko-sample-persistence-dc-java/pom.xml
@@ -67,9 +67,9 @@
<version>${pekko-persistence-cassandra.version}</version>
</dependency>
<dependency>
- <groupId>org.apache.pekko</groupId>
-
<artifactId>pekko-persistence-cassandra-launcher_${scala.binary.version}</artifactId>
- <version>${pekko-persistence-cassandra.version}</version>
+ <groupId>org.testcontainers</groupId>
+ <artifactId>testcontainers-cassandra</artifactId>
+ <version>2.0.5</version>
</dependency>
<dependency>
<groupId>org.apache.pekko</groupId>
diff --git
a/pekko-sample-persistence-dc-java/src/main/java/sample/persistence/res/MainApp.java
b/pekko-sample-persistence-dc-java/src/main/java/sample/persistence/res/MainApp.java
index 4577f15..22f2deb 100644
---
a/pekko-sample-persistence-dc-java/src/main/java/sample/persistence/res/MainApp.java
+++
b/pekko-sample-persistence-dc-java/src/main/java/sample/persistence/res/MainApp.java
@@ -1,6 +1,5 @@
package sample.persistence.res;
-import java.io.File;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
@@ -14,8 +13,9 @@ import
org.apache.pekko.cluster.sharding.typed.ReplicatedSharding;
import org.apache.pekko.cluster.sharding.typed.ReplicatedShardingExtension;
import org.apache.pekko.cluster.typed.Cluster;
import org.apache.pekko.management.javadsl.PekkoManagement;
-import org.apache.pekko.persistence.cassandra.testkit.CassandraLauncher;
import org.apache.pekko.persistence.typed.ReplicaId;
+import org.testcontainers.cassandra.CassandraContainer;
+import org.testcontainers.utility.DockerImageName;
import com.typesafe.config.Config;
import com.typesafe.config.ConfigFactory;
import sample.persistence.res.counter.ThumbsUpCounter;
@@ -81,12 +81,14 @@ public class MainApp {
* in a real application a pre-existing Cassandra cluster should be used.
*/
private static void startCassandraDatabase() {
- File databaseDirectory = new File("target/cassandra-db");
- CassandraLauncher.start(
- databaseDirectory,
- CassandraLauncher.DefaultTestConfigResource(),
- false,
- 9042);
+ final CassandraContainer container = new CassandraContainer(
+ DockerImageName.parse("cassandra:5.0.5"));
+ // bind to the fixed port 9042 so that the sample nodes can connect with
the default driver settings
+
container.setPortBindings(java.util.Collections.singletonList("9042:9042"));
+ container.start();
+
+ // shut the cassandra instance down when the JVM stops
+ Runtime.getRuntime().addShutdownHook(new Thread(() -> container.stop()));
}
}
diff --git
a/pekko-sample-persistence-dc-java/src/main/resources/application.conf
b/pekko-sample-persistence-dc-java/src/main/resources/application.conf
index 338edd0..badf454 100644
--- a/pekko-sample-persistence-dc-java/src/main/resources/application.conf
+++ b/pekko-sample-persistence-dc-java/src/main/resources/application.conf
@@ -50,6 +50,8 @@ pekko.persistence.cassandra {
}
datastax-java-driver.advanced.reconnect-on-init = true
+# creating the keyspace and tables in a freshly started Cassandra container
can exceed the default 2s timeout
+datastax-java-driver.basic.request.timeout = 10s
# Apache Pekko Management config:
https://pekko.apache.org/docs/pekko-management/current/index.html
diff --git a/pekko-sample-persistence-dc-scala/README.md
b/pekko-sample-persistence-dc-scala/README.md
index 097594b..665e012 100644
--- a/pekko-sample-persistence-dc-scala/README.md
+++ b/pekko-sample-persistence-dc-scala/README.md
@@ -6,6 +6,8 @@ to run a replica per datacenter.
## How to run
+Starting Apache Cassandra with the `cassandra` argument runs it in a Docker
container using [Testcontainers](https://testcontainers.com/), so Docker must
be installed and running.
+
1. In terminal 1: `sbt "runMain sample.persistence.res.MainApp cassandra"`
1. In terminal 2: `sbt "runMain sample.persistence.res.MainApp 7345 eu-west"`
diff --git a/pekko-sample-persistence-dc-scala/build.sbt
b/pekko-sample-persistence-dc-scala/build.sbt
index d74cd23..fe704cc 100644
--- a/pekko-sample-persistence-dc-scala/build.sbt
+++ b/pekko-sample-persistence-dc-scala/build.sbt
@@ -20,16 +20,11 @@ libraryDependencies ++= Seq(
"org.apache.pekko" %% "pekko-management" % pekkoClusterManagementVersion,
"org.apache.pekko" %% "pekko-management-cluster-http" %
pekkoClusterManagementVersion,
"org.apache.pekko" %% "pekko-persistence-cassandra" % cassandraPluginVersion,
- "org.apache.pekko" %% "pekko-persistence-cassandra-launcher" %
cassandraPluginVersion,
+ "org.testcontainers" % "testcontainers-cassandra" % "2.0.5",
"ch.qos.logback" % "logback-classic" % logbackVersion,
"org.apache.pekko" %% "pekko-persistence-testkit" % pekkoVersion % Test,
"org.scalatest" %% "scalatest" % "3.2.19" % Test)
-// transitive dependency of akka 2.5x that is brought in by addons but evicted
-dependencyOverrides += "org.apache.pekko" %% "pekko-protobuf" % pekkoVersion
-dependencyOverrides += "org.apache.pekko" %% "pekko-cluster-tools" %
pekkoVersion
-dependencyOverrides += "org.apache.pekko" %% "pekko-coordination" %
pekkoVersion
-
licenses := Seq(("CC0",
url("http://creativecommons.org/publicdomain/zero/1.0")))
// Startup aliases for the cassandra server and for two seed nodes one for
eu-west and another for eu-central
diff --git
a/pekko-sample-persistence-dc-scala/src/main/resources/application.conf
b/pekko-sample-persistence-dc-scala/src/main/resources/application.conf
index 248248d..fa90d16 100644
--- a/pekko-sample-persistence-dc-scala/src/main/resources/application.conf
+++ b/pekko-sample-persistence-dc-scala/src/main/resources/application.conf
@@ -50,6 +50,8 @@ pekko.persistence.cassandra {
}
datastax-java-driver.advanced.reconnect-on-init = true
+# creating the keyspace and tables in a freshly started Cassandra container
can exceed the default 2s timeout
+datastax-java-driver.basic.request.timeout = 10s
# Apache Pekko Management config:
https://pekko.apache.org/docs/pekko-management/current/index.html
diff --git
a/pekko-sample-persistence-dc-scala/src/main/scala/sample/persistence/res/MainApp.scala
b/pekko-sample-persistence-dc-scala/src/main/scala/sample/persistence/res/MainApp.scala
index cd78041..34d7966 100644
---
a/pekko-sample-persistence-dc-scala/src/main/scala/sample/persistence/res/MainApp.scala
+++
b/pekko-sample-persistence-dc-scala/src/main/scala/sample/persistence/res/MainApp.scala
@@ -1,6 +1,5 @@
package sample.persistence.res
-import java.io.File
import java.util.concurrent.CountDownLatch
import org.apache.pekko.actor.typed.ActorSystem
@@ -8,7 +7,6 @@ import org.apache.pekko.actor.typed.scaladsl.Behaviors
import org.apache.pekko.cluster.sharding.typed.{ ReplicatedSharding,
ReplicatedShardingExtension }
import org.apache.pekko.http.scaladsl.Http
import org.apache.pekko.management.scaladsl.PekkoManagement
-import org.apache.pekko.persistence.cassandra.testkit.CassandraLauncher
import org.apache.pekko.persistence.typed.ReplicaId
import com.typesafe.config.{ Config, ConfigFactory }
import sample.persistence.res.bank.BankAccount
@@ -88,12 +86,17 @@ object MainApp {
* in a real application a pre-existing Cassandra cluster should be used.
*/
def startCassandraDatabase(): Unit = {
- val databaseDirectory = new File("target/cassandra-db")
- CassandraLauncher.start(
- databaseDirectory,
- CassandraLauncher.DefaultTestConfigResource,
- clean = false,
- port = 9042)
+ import org.testcontainers.cassandra.CassandraContainer
+ import org.testcontainers.utility.DockerImageName
+ val container = new
CassandraContainer(DockerImageName.parse("cassandra:5.0.5"))
+ // bind to the fixed port 9042 so that the sample nodes can connect with
the default driver settings
+ container.setPortBindings(java.util.Collections.singletonList("9042:9042"))
+ container.start()
+
+ // shut the cassandra instance down when the JVM stops
+ sys.addShutdownHook {
+ container.stop()
+ }
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]