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]

Reply via email to