This is an automated email from the ASF dual-hosted git repository.
fanningpj pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-pekko-samples.git
The following commit(s) were added to refs/heads/main by this push:
new d66906e Created WorkResultConsumerActor for handling published events
(#67)
d66906e is described below
commit d66906ecbd6271ef7b30c5341cae49ee7547b3ee
Author: Mario Renau - Alstom
<[email protected]>
AuthorDate: Tue Aug 15 22:54:57 2023 +0200
Created WorkResultConsumerActor for handling published events (#67)
* Created WorkResultConsumerActor for handling published events
* reorder imports
---------
Co-authored-by: mario <[email protected]>
Co-authored-by: PJ Fanning <[email protected]>
---
.../src/main/scala/worker/Main.scala | 14 +++++----
.../src/main/scala/worker/WorkManager.scala | 3 ++
.../scala/worker/WorkResultConsumerActor.scala | 33 ++++++++++++++++++++++
3 files changed, 45 insertions(+), 5 deletions(-)
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 42a9558..03002fd 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
@@ -2,14 +2,12 @@ package worker
import java.io.File
import java.util.concurrent.CountDownLatch
-
import org.apache.pekko.actor.typed.ActorSystem
-import org.apache.pekko.actor.typed.scaladsl.Behaviors
-import org.apache.pekko.cluster.typed.Cluster
+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 }
-import org.apache.pekko.cluster.typed.SelfUp
-import org.apache.pekko.cluster.typed.Subscribe
object Main {
@@ -39,6 +37,11 @@ object Main {
}
}
+ private def createWorkResultsConsumer(ctx: ActorContext[SelfUp]): Unit = {
+ val eventListener = ctx.spawn(WorkResultConsumerActor(),
"work-consumer-actor")
+ ctx.system.eventStream ! EventStream.Subscribe(eventListener)
+ }
+
def startClusterInSameJvm(): Unit = {
startCassandraDatabase()
// two backend nodes
@@ -57,6 +60,7 @@ object Main {
Behaviors.setup[SelfUp](ctx => {
val cluster = Cluster(ctx.system)
cluster.subscriptions ! Subscribe(ctx.self, classOf[SelfUp])
+ createWorkResultsConsumer(ctx)
Behaviors.receiveMessage {
case SelfUp(_) =>
ctx.log.info("Node is up")
diff --git
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
index a2a07f0..11a2611 100644
---
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
+++
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
@@ -119,6 +119,9 @@ object WorkManager {
// Any in progress work from the previous incarnation is retried
ctx.self ! ResetWorkInProgress
}
+ // Publish events to the system event stream as PublishedEvent after
they have been persisted
+ .withEventPublishing(enabled = true)
+
}
}
diff --git
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
new file mode 100644
index 0000000..dbaf878
--- /dev/null
+++
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
@@ -0,0 +1,33 @@
+package worker
+
+import org.apache.pekko.actor.typed.Behavior
+import org.apache.pekko.actor.typed.scaladsl.{ ActorContext, Behaviors }
+import org.apache.pekko.persistence.typed.PublishedEvent
+import worker.WorkState.{ WorkAccepted, WorkCompleted, WorkInProgressReset,
WorkStarted }
+
+object WorkResultConsumerActor {
+ def apply(): Behavior[PublishedEvent] =
+ Behaviors.setup { context =>
+ context.log.info("WorkResultConsumerActor started")
+
+ Behaviors.receiveMessage { message =>
+ handleReceivedEvent(context, message)
+ Behaviors.same
+ }
+ }
+
+ private def handleReceivedEvent(context: ActorContext[PublishedEvent],
event: PublishedEvent): Unit = {
+ val actualEvent = event.event
+ event.event match {
+ case WorkInProgressReset =>
+ context.log.info(s"Received published event
[${actualEvent.getClass.getCanonicalName}]: {}", event)
+ case WorkCompleted(workId) =>
+ context.log.info(s"Received published event
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+ case WorkStarted(workId) =>
+ context.log.info(s"Received published event
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+ case WorkAccepted(workId) =>
+ context.log.info(s"Received published event
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+ case _ => context.log.warn("Message not supported")
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]