Skip to content
Merged
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
39 changes: 3 additions & 36 deletions benchmarks/ce/src/main/scala/sage/benchmarks/CeBenchmarks.scala
Original file line number Diff line number Diff line change
@@ -1,46 +1,13 @@
package sage.benchmarks

import java.util.concurrent.TimeUnit
import org.openjdk.jmh.annotations.Param

import org.openjdk.jmh.annotations.*

@State(Scope.Benchmark)
@BenchmarkMode(Array(Mode.Throughput))
@OutputTimeUnit(TimeUnit.SECONDS)
@OperationsPerInvocation(1000) // = Payloads.KeyCount
@Fork(1)
@Warmup(iterations = 3, time = 3)
@Measurement(iterations = 5, time = 3)
class ThroughputBench extends RedisBenchState {
class ThroughputBench extends ThroughputBenchBase {

@Param(Array("sage-ce", "redis4cats")) var client: String = "sage-ce"
@Param(Array("1", "8", "64", "256")) var concurrency: Int = 1
@Param(Array("16", "1024")) var valueSize: Int = 16

protected def subjectName: String = client
override protected def seedValueBytes: Int = valueSize
protected def buildClient(host: String, port: Int, name: String): BenchClient = Clients.build(host, port, name)

@Benchmark def get(): Long = subject.getAll(keys, concurrency)
@Benchmark def set(): Long = subject.setAll(keys, Payloads.value(valueSize), concurrency)
}

// Run one large-reply command per invocation. The throughput benchmark covers concurrent requests. valueSize controls the seeded value size.
@State(Scope.Benchmark)
@BenchmarkMode(Array(Mode.Throughput))
@OutputTimeUnit(TimeUnit.SECONDS)
@Fork(1)
@Warmup(iterations = 3, time = 3)
@Measurement(iterations = 5, time = 3)
class CollectionBench extends RedisBenchState {
class CollectionBench extends CollectionBenchBase {

@Param(Array("sage-ce", "redis4cats")) var client: String = "sage-ce"
@Param(Array("16")) var valueSize: Int = 16

protected def subjectName: String = client
override protected def seedValueBytes: Int = valueSize
protected def buildClient(host: String, port: Int, name: String): BenchClient = Clients.build(host, port, name)

@Benchmark def mget(): Long = subject.mget(keys)
@Benchmark def hgetall(): Long = subject.hgetall(Payloads.HashKey)
}
68 changes: 10 additions & 58 deletions benchmarks/ce/src/main/scala/sage/benchmarks/Clients.scala
Original file line number Diff line number Diff line change
Expand Up @@ -21,77 +21,29 @@ object Clients {
}
}

final class SageCeBench(host: String, port: Int) extends BenchClient {
final class SageCeBench(host: String, port: Int) extends SageBench[IO] {

private val client: SageClient =
protected val client: SageClient =
SageClient.connect(SageConfig(topology = Topology.Standalone(Endpoint(host, port)))).unsafeRunSync()

def name: String = "sage-ce"
protected def run[A](effect: IO[A]): Unit = effect.unsafeRunSync(): Unit

def seed(prefix: String, count: Int, value: String, hashKey: String, fields: Int): Unit = {
val sets = (0 until count).toList.traverse_(i => client.set(s"$prefix:$i", value))
val hash = (0 until fields).map(i => (s"f$i", value)).toList match {
case h :: t => client.hSet(hashKey, h, t*).void
case Nil => IO.unit
}
(sets *> hash).unsafeRunSync()
}

def getAll(keys: Array[String], concurrency: Int): Long =
Payloads
.groups(keys, concurrency)
.toList
.parTraverse(_.toList.traverse(client.get[String]))
.map(_.flatten.flatten.map(_.length.toLong).sum)
.unsafeRunSync()

def setAll(keys: Array[String], value: String, concurrency: Int): Long =
Payloads
.groups(keys, concurrency)
.toList
.parTraverse_(_.toList.traverse_(client.set(_, value)))
.as(keys.length.toLong)
.unsafeRunSync()

def mget(keys: Array[String]): Long =
client.mGet[String](keys.head, keys.tail*).map(_.flatten.map(_.length.toLong).sum).unsafeRunSync()

def hgetall(key: String): Long = client.hGetAll[String, String](key).map(_.size.toLong).unsafeRunSync()

def close(): Unit = client.close.unsafeRunSync()
protected def inLanes[A](work: Payloads.Workload)(perKey: String => IO[A]): IO[Unit] = work.lanes.parTraverse_(_.traverse_(perKey))
}

final class Redis4catsBench(host: String, port: Int) extends BenchClient {

private val (redis, release) = Redis[IO].utf8(s"redis://$host:$port").allocated.unsafeRunSync()

def name: String = "redis4cats"

def seed(prefix: String, count: Int, value: String, hashKey: String, fields: Int): Unit = {
val sets = (0 until count).toList.traverse_(i => redis.set(s"$prefix:$i", value))
val hash = (0 until fields).toList.traverse_(i => redis.hSet(hashKey, s"f$i", value))
(sets *> hash).unsafeRunSync()
}

def getAll(keys: Array[String], concurrency: Int): Long =
Payloads
.groups(keys, concurrency)
.toList
.parTraverse(_.toList.traverse(redis.get))
.map(_.flatten.flatten.map(_.length.toLong).sum)
.unsafeRunSync()
def getAll(work: Payloads.Workload): Unit =
work.lanes.parTraverse_(_.traverse_(redis.get)).unsafeRunSync()

def setAll(keys: Array[String], value: String, concurrency: Int): Long =
Payloads
.groups(keys, concurrency)
.toList
.parTraverse_(_.toList.traverse_(redis.set(_, value)))
.as(keys.length.toLong)
.unsafeRunSync()
def setAll(work: Payloads.Workload, value: String): Unit =
work.lanes.parTraverse_(_.traverse_(redis.set(_, value))).unsafeRunSync()

def mget(keys: Array[String]): Long = redis.mGet(keys.toSet).map(_.values.map(_.length.toLong).sum).unsafeRunSync()
def mget(): Unit = redis.mGet(Payloads.Keys.set).void.unsafeRunSync()

def hgetall(key: String): Long = redis.hGetAll(key).map(_.size.toLong).unsafeRunSync()
def hgetall(): Unit = redis.hGetAll(Payloads.HashKey).void.unsafeRunSync()

def close(): Unit = release.unsafeRunSync()
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
package sage.benchmarks

object Clients {
def build(host: String, port: Int, name: String): BenchClient = throw new IllegalArgumentException(s"unknown client: $name")
}
34 changes: 5 additions & 29 deletions benchmarks/kyo/src/main/scala/sage/benchmarks/Clients.scala
Original file line number Diff line number Diff line change
Expand Up @@ -22,37 +22,13 @@ private object Run {
KyoApp.Unsafe.runAndBlock(Duration.Infinity)(program).getOrThrow
}

final class SageKyoBench(host: String, port: Int) extends BenchClient {
final class SageKyoBench(host: String, port: Int) extends SageBench[[A] =>> A < (Abort[SageException] & Async)] {

private val client: SageClient =
protected val client: SageClient =
Run(SageClient.connect(SageConfig(topology = Topology.Standalone(Endpoint(host, port)))))

def name: String = "sage-kyo"
protected def run[A](effect: A < (Abort[SageException] & Async)): Unit = Run(effect): Unit

def seed(prefix: String, count: Int, value: String, hashKey: String, fields: Int): Unit =
Run(for {
_ <- Kyo.foreachDiscard(0 until count)(i => client.set(s"$prefix:$i", value))
_ <- Kyo.foreachDiscard(0 until fields)(i => client.hSet(hashKey, (s"f$i", value)))
} yield ())

def getAll(keys: Array[String], concurrency: Int): Long =
Run(
Async
.foreach(Payloads.groups(keys, concurrency).toList, concurrency)(g => Kyo.foreach(g.toList)(k => client.get[String](k)).map(_.toList))
.map(_.toList.flatten.flatten.map(_.length.toLong).sum)
)

def setAll(keys: Array[String], value: String, concurrency: Int): Long =
Run(
Async
.foreachDiscard(Payloads.groups(keys, concurrency).toList, concurrency)(g => Kyo.foreachDiscard(g.toList)(k => client.set(k, value)))
.map(_ => keys.length.toLong)
)

def mget(keys: Array[String]): Long =
Run(client.mGet[String](keys.head, keys.tail*).map(_.flatten.map(_.length.toLong).sum))

def hgetall(key: String): Long = Run(client.hGetAll[String, String](key).map(_.size.toLong))

def close(): Unit = Run(client.close)
protected def inLanes[A](work: Payloads.Workload)(perKey: String => A < (Abort[SageException] & Async)): Unit < (Abort[SageException] & Async) =
Async.foreachDiscard(work.lanes, work.concurrency)(Kyo.foreachDiscard(_)(perKey))
}
41 changes: 4 additions & 37 deletions benchmarks/kyo/src/main/scala/sage/benchmarks/KyoBenchmarks.scala
Original file line number Diff line number Diff line change
@@ -1,46 +1,13 @@
package sage.benchmarks

import java.util.concurrent.TimeUnit
import org.openjdk.jmh.annotations.Param

import org.openjdk.jmh.annotations.*
class ThroughputBench extends ThroughputBenchBase {

@State(Scope.Benchmark)
@BenchmarkMode(Array(Mode.Throughput))
@OutputTimeUnit(TimeUnit.SECONDS)
@OperationsPerInvocation(1000) // = Payloads.KeyCount
@Fork(1)
@Warmup(iterations = 3, time = 3)
@Measurement(iterations = 5, time = 3)
class ThroughputBench extends RedisBenchState {

@Param(Array("sage-kyo")) var client: String = "sage-kyo"
@Param(Array("1", "8", "64", "256")) var concurrency: Int = 1
@Param(Array("16", "1024")) var valueSize: Int = 16

protected def subjectName: String = client
override protected def seedValueBytes: Int = valueSize
protected def buildClient(host: String, port: Int, name: String): BenchClient = Clients.build(host, port, name)

@Benchmark def get(): Long = subject.getAll(keys, concurrency)
@Benchmark def set(): Long = subject.setAll(keys, Payloads.value(valueSize), concurrency)
@Param(Array("sage-kyo")) var client: String = "sage-kyo"
}

// Run one large-reply command per invocation. The throughput benchmark covers concurrent requests. valueSize controls the seeded value size.
@State(Scope.Benchmark)
@BenchmarkMode(Array(Mode.Throughput))
@OutputTimeUnit(TimeUnit.SECONDS)
@Fork(1)
@Warmup(iterations = 3, time = 3)
@Measurement(iterations = 5, time = 3)
class CollectionBench extends RedisBenchState {
class CollectionBench extends CollectionBenchBase {

@Param(Array("sage-kyo")) var client: String = "sage-kyo"
@Param(Array("16")) var valueSize: Int = 16

protected def subjectName: String = client
override protected def seedValueBytes: Int = valueSize
protected def buildClient(host: String, port: Int, name: String): BenchClient = Clients.build(host, port, name)

@Benchmark def mget(): Long = subject.mget(keys)
@Benchmark def hgetall(): Long = subject.hgetall(Payloads.HashKey)
}
Loading