This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-persistence-dynamodb.git
The following commit(s) were added to refs/heads/main by this push:
new c1aa76e refactor: replace deprecated `extends App` with explicit main
method (#334)
c1aa76e is described below
commit c1aa76eb8339830c34d02e8705e38b1e408ae777
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Mon Jul 6 16:26:56 2026 +0800
refactor: replace deprecated `extends App` with explicit main method (#334)
Motivation:
`extends App` is deprecated in Scala 3. Replace with explicit `def main`
method for forward compatibility.
Modification:
Replace `object X extends App with Trait { ... }` with `object X extends
Trait { def main(args: Array[String]): Unit = { ... } }` in
WriteThroughputBench.
Result:
No more usage of deprecated `scala.App` trait. Code is compatible with
Scala 3.
Tests:
Not run - docs only
References:
None - Scala 3 compatibility
---
.../dynamodb/journal/WriteThroughputBench.scala | 258 +++++++++++----------
1 file changed, 130 insertions(+), 128 deletions(-)
diff --git
a/src/test/scala-2/org/apache/pekko/persistence/dynamodb/journal/WriteThroughputBench.scala
b/src/test/scala-2/org/apache/pekko/persistence/dynamodb/journal/WriteThroughputBench.scala
index 8a194eb..7023210 100644
---
a/src/test/scala-2/org/apache/pekko/persistence/dynamodb/journal/WriteThroughputBench.scala
+++
b/src/test/scala-2/org/apache/pekko/persistence/dynamodb/journal/WriteThroughputBench.scala
@@ -26,74 +26,75 @@ import java.util.UUID
import java.util.concurrent.ThreadLocalRandom
import scala.concurrent.duration._
-object WriteThroughputBench extends App with DynamoDBUtils {
-
- def rnd = ThreadLocalRandom.current()
- final val oneBillion = 1000L * 1000 * 1000
-
- class H private (private val entries: Map[Int, Int]) {
- def this(i: Int) = this(Map(i -> 1))
- def this() = this(Map.empty[Int, Int])
- def +(other: H): H = {
- val merged =
- other.entries.foldLeft(entries)((acc, pair) => acc.updated(pair._1,
pair._2 + entries.getOrElse(pair._1, 0)))
- new H(merged)
+object WriteThroughputBench extends DynamoDBUtils {
+ def main(args: Array[String]): Unit = {
+
+ def rnd = ThreadLocalRandom.current()
+ final val oneBillion = 1000L * 1000 * 1000
+
+ class H private (private val entries: Map[Int, Int]) {
+ def this(i: Int) = this(Map(i -> 1))
+ def this() = this(Map.empty[Int, Int])
+ def +(other: H): H = {
+ val merged =
+ other.entries.foldLeft(entries)((acc, pair) => acc.updated(pair._1,
pair._2 + entries.getOrElse(pair._1, 0)))
+ new H(merged)
+ }
+ def record(value: Int): H = new H(entries.updated(value,
entries.getOrElse(value, 0) + 1))
+ override def toString: String =
+ if (entries.nonEmpty) {
+ val max = entries.keys.max
+ (0 to max).map(entries.getOrElse(_, 0)).mkString("[", ",", "]")
+ } else "empty"
}
- def record(value: Int): H = new H(entries.updated(value,
entries.getOrElse(value, 0) + 1))
- override def toString: String =
- if (entries.nonEmpty) {
- val max = entries.keys.max
- (0 to max).map(entries.getOrElse(_, 0)).mkString("[", ",", "]")
- } else "empty"
- }
- case object Go
- case class Report(endToEnd: Histogram, calls: Histogram, retries: H) {
- def +(other: Report): Report = {
- endToEnd.add(other.endToEnd)
- calls.add(other.calls)
- Report(endToEnd, calls, retries + other.retries)
+ case object Go
+ case class Report(endToEnd: Histogram, calls: Histogram, retries: H) {
+ def +(other: Report): Report = {
+ endToEnd.add(other.endToEnd)
+ calls.add(other.calls)
+ Report(endToEnd, calls, retries + other.retries)
+ }
+ }
+ object Report {
+ def apply(): Report = Report(new Histogram(3), new Histogram(3), new H)
+ def endToEnd(h: Histogram): Report = Report(h, new Histogram(3), new H)
+ def calls(h: Histogram, r: H): Report = Report(new Histogram(3), h, r)
}
- }
- object Report {
- def apply(): Report = Report(new Histogram(3), new Histogram(3), new H)
- def endToEnd(h: Histogram): Report = Report(h, new Histogram(3), new H)
- def calls(h: Histogram, r: H): Report = Report(new Histogram(3), h, r)
- }
- class Writer(stats: ActorRef) extends PersistentActor {
- val event = new Array[Byte](100)
- rnd.nextBytes(event)
-
- var nextReport = System.nanoTime + oneBillion
- val histo = new Histogram(3)
- val persistenceId = UUID.randomUUID().toString
-
- self ! Go
-
- def receiveCommand = {
- case Go =>
- val start = System.nanoTime
- persist(event) { _ =>
- val end = System.nanoTime
- histo.recordValue(end - start)
- if (end > nextReport) {
- stats ! Report.endToEnd(histo.copy())
- histo.reset()
- nextReport += oneBillion
+ class Writer(stats: ActorRef) extends PersistentActor {
+ val event = new Array[Byte](100)
+ rnd.nextBytes(event)
+
+ var nextReport = System.nanoTime + oneBillion
+ val histo = new Histogram(3)
+ val persistenceId = UUID.randomUUID().toString
+
+ self ! Go
+
+ def receiveCommand = {
+ case Go =>
+ val start = System.nanoTime
+ persist(event) { _ =>
+ val end = System.nanoTime
+ histo.recordValue(end - start)
+ if (end > nextReport) {
+ stats ! Report.endToEnd(histo.copy())
+ histo.reset()
+ nextReport += oneBillion
+ }
+ self ! Go
}
- self ! Go
- }
+ }
+ def receiveRecover = Actor.emptyBehavior
}
- def receiveRecover = Actor.emptyBehavior
- }
- val config =
- ConfigFactory
- .systemProperties()
- .withFallback(
- ConfigFactory
- .parseString("""
+ val config =
+ ConfigFactory
+ .systemProperties()
+ .withFallback(
+ ConfigFactory
+ .parseString("""
my-dynamodb-journal {
journal-table = "WriteThroughputBench"
endpoint = ${?AWS_DYNAMODB_ENDPOINT}
@@ -125,79 +126,80 @@ writer-dispatcher {
}
}
""").resolve)
- .withFallback(ConfigFactory.load())
+ .withFallback(ConfigFactory.load())
- implicit val system: ActorSystem = ActorSystem("WriteThroughputBench",
config)
- implicit val materializer: Materializer = Materializer(system)
+ implicit val system: ActorSystem = ActorSystem("WriteThroughputBench",
config)
+ implicit val materializer: Materializer = Materializer(system)
- /*
- * You will want to make sure that the table is deployed with the proper
values for
- * read and write throughput; the default is 10/10 which is incredibly low,
but defaulting
- * to larger values can burn through a big budget very quickly.
- */
- ensureJournalTableExists()
+ /*
+ * You will want to make sure that the table is deployed with the proper
values for
+ * read and write throughput; the default is 10/10 which is incredibly
low, but defaulting
+ * to larger values can burn through a big budget very quickly.
+ */
+ ensureJournalTableExists()
- val writers = system.settings.config.getInt("writers")
+ val writers = system.settings.config.getInt("writers")
- val completionMatcher: PartialFunction[Any, CompletionStrategy] = {
- case pekko.actor.Status.Success(s: CompletionStrategy) => s
- case pekko.actor.Status.Success(_) =>
CompletionStrategy.Draining
- case pekko.actor.Status.Success =>
CompletionStrategy.Draining
- }
+ val completionMatcher: PartialFunction[Any, CompletionStrategy] = {
+ case pekko.actor.Status.Success(s: CompletionStrategy) => s
+ case pekko.actor.Status.Success(_) =>
CompletionStrategy.Draining
+ case pekko.actor.Status.Success =>
CompletionStrategy.Draining
+ }
- val failureMatcher: PartialFunction[Any, Throwable] = {
- case pekko.actor.Status.Failure(cause) => cause
- }
+ val failureMatcher: PartialFunction[Any, Throwable] = {
+ case pekko.actor.Status.Failure(cause) => cause
+ }
- val endToEnd = Source
- .actorRef[Report](completionMatcher, failureMatcher, 3 * writers,
OverflowStrategy.dropHead)
- .conflate(_ + _)
- .prepend(Source.single(Report()))
- .expand(Iterator.continually(_))
- .withAttributes(Attributes.asyncBoundary)
-
- val calls = Source
- .actorRef[LatencyReport](completionMatcher, failureMatcher, 1000,
OverflowStrategy.dropHead)
- .conflateWithSeed(r => ({ val h = new Histogram(3);
h.recordValue(r.nanos); h }, new H(r.retries))) {
- case ((hist, h), LatencyReport(nanos, retries)) =>
- hist.recordValue(nanos)
- (hist, h.record(retries))
+ val endToEnd = Source
+ .actorRef[Report](completionMatcher, failureMatcher, 3 * writers,
OverflowStrategy.dropHead)
+ .conflate(_ + _)
+ .prepend(Source.single(Report()))
+ .expand(Iterator.continually(_))
+ .withAttributes(Attributes.asyncBoundary)
+
+ val calls = Source
+ .actorRef[LatencyReport](completionMatcher, failureMatcher, 1000,
OverflowStrategy.dropHead)
+ .conflateWithSeed(r => ({ val h = new Histogram(3);
h.recordValue(r.nanos); h }, new H(r.retries))) {
+ case ((hist, h), LatencyReport(nanos, retries)) =>
+ hist.recordValue(nanos)
+ (hist, h.record(retries))
+ }
+ .map(p => Report.calls(p._1, p._2))
+ .prepend(Source.single(Report()))
+ .expand(Iterator.continually(_))
+ .withAttributes(Attributes.asyncBoundary)
+
+ val (eRef, cRef) =
+ RunnableGraph
+ .fromGraph(GraphDSL.createGraph(endToEnd, calls)(Keep.both) { implicit
b => (e, c) =>
+ val zip = b.add(ZipWith((_: Unit, er: Report, cr: Report) => er +
cr))
+ Source.tick(1.second, 1.second, ()) ~> zip.in0
+ e ~> zip.in1
+ c ~> zip.in2
+ zip.out ~> Sink.foreach(printStats)
+ ClosedShape
+ })
+ .withAttributes(Attributes.inputBuffer(1, 1))
+ .run()
+
+ println(s"starting $writers writers")
+ val props = Props(new Writer(eRef)).withDispatcher("writer-dispatcher")
+ val writerRefs = (1 to writers).map(_ => system.actorOf(props))
+
+ Persistence(system).journalFor("") ! SetDBHelperReporter(cRef)
+
+ println("press <enter> to stop")
+ scala.io.StdIn.readLine()
+
+ system.terminate()
+ dynamo.shutdown()
+
+ def printStats(r: Report): Unit = {
+ def p(h: Histogram, pc: Double) = h.getValueAtPercentile(pc) / 1000000d
+ def perc(h: Histogram) =
+ f"50=${p(h, 0.5)}%5.1f 90=${p(h, 0.9)}%5.1f 99=${p(h, 0.99)}%5.1f
99.9=${p(h, 0.999)}%5.1f 99.99=${p(h, 0.9999)}%5.1f"
+ println(f"count ${r.endToEnd.getTotalCount}%6d/s percentiles:
endToEnd(${perc(r.endToEnd)}) calls(${perc(
+ r.calls)}) retries: ${r.retries}")
}
- .map(p => Report.calls(p._1, p._2))
- .prepend(Source.single(Report()))
- .expand(Iterator.continually(_))
- .withAttributes(Attributes.asyncBoundary)
-
- val (eRef, cRef) =
- RunnableGraph
- .fromGraph(GraphDSL.createGraph(endToEnd, calls)(Keep.both) { implicit b
=> (e, c) =>
- val zip = b.add(ZipWith((_: Unit, er: Report, cr: Report) => er + cr))
- Source.tick(1.second, 1.second, ()) ~> zip.in0
- e ~> zip.in1
- c ~> zip.in2
- zip.out ~> Sink.foreach(printStats)
- ClosedShape
- })
- .withAttributes(Attributes.inputBuffer(1, 1))
- .run()
-
- println(s"starting $writers writers")
- val props = Props(new Writer(eRef)).withDispatcher("writer-dispatcher")
- val writerRefs = (1 to writers).map(_ => system.actorOf(props))
-
- Persistence(system).journalFor("") ! SetDBHelperReporter(cRef)
-
- println("press <enter> to stop")
- scala.io.StdIn.readLine()
-
- system.terminate()
- dynamo.shutdown()
-
- def printStats(r: Report): Unit = {
- def p(h: Histogram, pc: Double) = h.getValueAtPercentile(pc) / 1000000d
- def perc(h: Histogram) =
- f"50=${p(h, 0.5)}%5.1f 90=${p(h, 0.9)}%5.1f 99=${p(h, 0.99)}%5.1f
99.9=${p(h, 0.999)}%5.1f 99.99=${p(h, 0.9999)}%5.1f"
- println(f"count ${r.endToEnd.getTotalCount}%6d/s percentiles:
endToEnd(${perc(r.endToEnd)}) calls(${perc(
- r.calls)}) retries: ${r.retries}")
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]