From 963469d0ea6f6cecfd6c328bbb9b3310fa2df8f0 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 18 Feb 2024 21:12:22 +0900 Subject: [PATCH 1/2] =?UTF-8?q?test:=20Netx=20client=20=EB=B6=80=ED=95=98?= =?UTF-8?q?=ED=85=8C=EC=8A=A4=ED=8A=B8=EB=A5=BC=20=EC=9E=91=EC=84=B1?= =?UTF-8?q?=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../redis/RedisStreamTransactionListener.kt | 9 +++ .../org/rooftop/netx/client/LoadRunner.kt | 20 ++++++ .../org/rooftop/netx/client/NetxClient.kt | 26 ++++++++ .../org/rooftop/netx/client/NetxTest.kt | 64 +++++++++++++++++++ .../netx/client/TransactionReceiveStorage.kt | 57 +++++++++++++++++ 5 files changed, 176 insertions(+) create mode 100644 src/test/kotlin/org/rooftop/netx/client/LoadRunner.kt create mode 100644 src/test/kotlin/org/rooftop/netx/client/NetxClient.kt create mode 100644 src/test/kotlin/org/rooftop/netx/client/NetxTest.kt create mode 100644 src/test/kotlin/org/rooftop/netx/client/TransactionReceiveStorage.kt diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt index 7252c94..b388eca 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt @@ -1,5 +1,6 @@ package org.rooftop.netx.redis +import io.lettuce.core.RedisBusyException import org.rooftop.netx.engine.AbstractTransactionDispatcher import org.rooftop.netx.engine.AbstractTransactionListener import org.rooftop.netx.idl.Transaction @@ -10,8 +11,10 @@ import org.springframework.data.redis.connection.stream.StreamOffset import org.springframework.data.redis.core.ReactiveRedisTemplate import org.springframework.data.redis.stream.StreamReceiver import reactor.core.publisher.Flux +import reactor.core.publisher.Mono import reactor.core.scheduler.Schedulers import kotlin.time.Duration.Companion.hours +import kotlin.time.Duration.Companion.milliseconds import kotlin.time.toJavaDuration class RedisStreamTransactionListener( @@ -42,6 +45,12 @@ class RedisStreamTransactionListener( private fun createGroupIfNotExists(transactionId: String): Flux { return reactiveRedisTemplate.opsForStream() .createGroup(transactionId, ReadOffset.from("0"), nodeGroup) + .onErrorResume { + if (it.cause is RedisBusyException) { + return@onErrorResume Mono.just(transactionId) + } + throw it + } .flatMapMany { Flux.just(it) } } } diff --git a/src/test/kotlin/org/rooftop/netx/client/LoadRunner.kt b/src/test/kotlin/org/rooftop/netx/client/LoadRunner.kt new file mode 100644 index 0000000..eaada4f --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/client/LoadRunner.kt @@ -0,0 +1,20 @@ +package org.rooftop.netx.client + +import org.springframework.boot.test.context.TestComponent +import reactor.core.publisher.Flux +import reactor.core.scheduler.Schedulers + +@TestComponent +class LoadRunner { + + fun load(count: Int, behavior: Runnable) { + val iter = mutableListOf() + for (i in 1..count) { + iter.add(behavior) + } + Flux.fromIterable(iter) + .publishOn(Schedulers.boundedElastic()) + .map { it.run() } + .subscribe() + } +} diff --git a/src/test/kotlin/org/rooftop/netx/client/NetxClient.kt b/src/test/kotlin/org/rooftop/netx/client/NetxClient.kt new file mode 100644 index 0000000..691875a --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/client/NetxClient.kt @@ -0,0 +1,26 @@ +package org.rooftop.netx.client + +import org.rooftop.netx.api.TransactionManager +import org.springframework.stereotype.Service + +@Service +class NetxClient( + private val transactionManager: TransactionManager, +) { + + fun startTransaction(undo: String): String { + return transactionManager.start(undo).block()!! + } + + fun rollbackTransaction(transactionId: String, cause: String): String { + return transactionManager.rollback(transactionId, cause).block()!! + } + + fun joinTransaction(transactionId: String, undo: String): String { + return transactionManager.join(transactionId, undo).block()!! + } + + fun commitTransaction(transactionId: String): String { + return transactionManager.commit(transactionId).block()!! + } +} diff --git a/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt b/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt new file mode 100644 index 0000000..30501e1 --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt @@ -0,0 +1,64 @@ +package org.rooftop.netx.client + +import io.kotest.assertions.nondeterministic.eventually +import io.kotest.core.annotation.DisplayName +import io.kotest.core.spec.style.FunSpec +import io.kotest.data.forAll +import io.kotest.data.row +import org.rooftop.netx.meta.EnableDistributedTransaction +import org.rooftop.netx.redis.RedisContainer +import org.springframework.boot.test.context.SpringBootTest +import reactor.core.publisher.Hooks +import kotlin.time.Duration.Companion.minutes + +@DisplayName("Netx 테스트의") +@SpringBootTest( + classes = [ + RedisContainer::class, + LoadRunner::class, + NetxClient::class, + TransactionReceiveStorage::class, + ] +) +@EnableDistributedTransaction +internal class NetxTest( + private val netxClient: NetxClient, + private val loadRunner: LoadRunner, + private val transactionReceiveStorage: TransactionReceiveStorage, +) : FunSpec({ + + test("Netx는 부하가 가중되어도, 결과적 일관성을 보장한다.") { + forAll( + row(1, 1), + row(10, 10), + row(100, 100), + row(1_000, 1_000), + row(10_000, 10_000), + row(100_000, 100_000), + row(1_000_000, 1_000_000), + ) { commitLoadCount, rollbackLoadCount -> + transactionReceiveStorage.clear() + + loadRunner.load(commitLoadCount) { + val transactionId = netxClient.startTransaction("") + netxClient.joinTransaction(transactionId, "") + netxClient.commitTransaction(transactionId) + } + + loadRunner.load(rollbackLoadCount) { + val transactionId = netxClient.startTransaction("") + netxClient.joinTransaction(transactionId, "") + netxClient.rollbackTransaction(transactionId, "") + } + + eventually(10.minutes) { + transactionReceiveStorage.startCountShouldBe(commitLoadCount + rollbackLoadCount) + transactionReceiveStorage.joinCountShouldBe(commitLoadCount + rollbackLoadCount) + transactionReceiveStorage.commitCountShouldBe(commitLoadCount) + transactionReceiveStorage.rollbackCountShouldBe(rollbackLoadCount) + } + + Thread.sleep(Long.MAX_VALUE) + } + } +}) diff --git a/src/test/kotlin/org/rooftop/netx/client/TransactionReceiveStorage.kt b/src/test/kotlin/org/rooftop/netx/client/TransactionReceiveStorage.kt new file mode 100644 index 0000000..94eb067 --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/client/TransactionReceiveStorage.kt @@ -0,0 +1,57 @@ +package org.rooftop.netx.client + +import io.kotest.matchers.shouldBe +import org.rooftop.netx.api.* +import org.rooftop.netx.meta.TransactionHandler +import reactor.core.publisher.Mono + +@TransactionHandler +class TransactionReceiveStorage( + private val storage: MutableMap>, +) { + + fun clear() { + storage.clear() + } + + fun joinCountShouldBe(count: Int) { + (storage["JOIN"]?.size ?: 0) shouldBe count + } + + fun startCountShouldBe(count: Int) { + (storage["START"]?.size ?: 0) shouldBe count + } + + fun commitCountShouldBe(count: Int) { + (storage["COMMIT"]?.size ?: 0) shouldBe count + } + + fun rollbackCountShouldBe(count: Int) { + (storage["ROLLBACK"]?.size ?: 0) shouldBe count + } + + @TransactionRollbackHandler + fun logRollback(transaction: TransactionRollbackEvent): Mono { + return Mono.fromCallable { log("ROLLBACK", transaction) } + } + + @TransactionStartHandler + fun logStart(transaction: TransactionStartEvent): Mono { + return Mono.fromCallable { log("START", transaction) } + } + + @TransactionJoinHandler + fun logJoin(transaction: TransactionJoinEvent): Mono { + return Mono.fromCallable { log("JOIN", transaction) } + } + + @TransactionCommitHandler + fun logCommit(transaction: TransactionCommitEvent): Mono { + return Mono.fromCallable { log("COMMIT", transaction) } + } + + private fun log(key: String, transaction: TransactionEvent) { + storage.putIfAbsent(key, mutableListOf()) + storage[key]?.add(transaction) + } +} From c211f4ad6e5e3a267f9347b44dfb6cee62358294 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 18 Feb 2024 21:13:45 +0900 Subject: [PATCH 2/2] =?UTF-8?q?perf:=20netx=20redis-stream=20=EA=B5=AC?= =?UTF-8?q?=ED=98=84=EC=B1=84=EA=B0=80=20=ED=95=98=EB=82=98=EC=9D=98=20?= =?UTF-8?q?=EC=BB=A4=EB=84=A5=EC=85=98=EC=9C=BC=EB=A1=9C=20=ED=8A=B8?= =?UTF-8?q?=EB=9E=9C=EC=9E=AD=EC=85=98=EC=9D=84=20=EC=88=98=EC=8B=A0?= =?UTF-8?q?=ED=95=A0=20=EC=88=98=20=EC=9E=88=EB=8F=84=EB=A1=9D=20=EC=88=98?= =?UTF-8?q?=EC=A0=95=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../engine/AbstractTransactionDispatcher.kt | 6 -- .../engine/AbstractTransactionListener.kt | 9 ++- .../netx/engine/AbstractTransactionManager.kt | 39 ++++------- .../AbstractTransactionRetrySupporter.kt | 2 - .../redis/RedisStreamTransactionDispatcher.kt | 12 ++-- .../redis/RedisStreamTransactionListener.kt | 17 +++-- .../redis/RedisStreamTransactionManager.kt | 43 ++++++++---- .../redis/RedisStreamTransactionRemover.kt | 51 -------------- .../netx/redis/RedisTransactionConfigurer.kt | 22 ++----- .../redis/RedisTransactionRetrySupporter.kt | 25 ++++--- .../org/rooftop/netx/client/NetxTest.kt | 7 +- .../NoAckRedisStreamTransactionDispatcher.kt | 6 -- .../redis/NoAckRedisTransactionConfigurer.kt | 22 ++----- .../org/rooftop/netx/redis/RedisAssertions.kt | 2 +- .../RedisStreamTransactionManagerTest.kt | 1 - .../RedisStreamTransactionRemoverTest.kt | 66 ------------------- 16 files changed, 88 insertions(+), 242 deletions(-) delete mode 100644 src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemover.kt delete mode 100644 src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemoverTest.kt diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt index 4a2fc0e..3a351e6 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt @@ -17,7 +17,6 @@ abstract class AbstractTransactionDispatcher { fun dispatch(transaction: Transaction, messageId: String): Flux { return Mono.just(transaction.state) - .doOnNext { deleteElastic(transaction, messageId) } .filter { state -> transactionHandlerFunctions.containsKey(state) } .flatMapMany { state -> Flux.fromIterable( @@ -84,11 +83,6 @@ abstract class AbstractTransactionDispatcher { messageId: String ): Mono> - protected abstract fun deleteElastic( - transaction: Transaction, - messageId: String - ) - private companion object { private val cannotFindMatchedTransactionEventException = java.lang.IllegalStateException("Cannot find matched transaction event") diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionListener.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionListener.kt index 05e31e8..262645e 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionListener.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionListener.kt @@ -2,18 +2,21 @@ package org.rooftop.netx.engine import org.rooftop.netx.idl.Transaction import reactor.core.publisher.Flux +import reactor.core.scheduler.Schedulers abstract class AbstractTransactionListener( private val transactionDispatcher: AbstractTransactionDispatcher, ) { - fun subscribeStream(transactionId: String): Flux> { - return receive(transactionId) + fun subscribeStream() { + receive() .flatMap { (transaction, messageId) -> transactionDispatcher.dispatch(transaction, messageId) .map { transaction to messageId } } + .subscribeOn(Schedulers.parallel()) + .subscribe() } - protected abstract fun receive(transactionId: String): Flux> + protected abstract fun receive(): Flux> } diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index 38ea9ad..a8737de 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -5,21 +5,16 @@ import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState import org.rooftop.netx.idl.transaction import reactor.core.publisher.Mono -import reactor.core.scheduler.Schedulers abstract class AbstractTransactionManager( nodeId: Int, private val nodeGroup: String, private val nodeName: String, private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId), - private val transactionListener: AbstractTransactionListener, - private val transactionRetrySupporter: AbstractTransactionRetrySupporter, ) : TransactionManager { final override fun start(undo: String): Mono { return startTransaction(undo) - .subscribeTransaction() - .watchTransaction() .contextWrite { it.put(CONTEXT_TX_KEY, transactionIdGenerator.generate()) } } @@ -37,10 +32,16 @@ abstract class AbstractTransactionManager( } final override fun join(transactionId: String, undo: String): Mono { - return exists(transactionId) + return findAnyTransaction(transactionId) + .map { + if (it == TransactionState.TRANSACTION_STATE_ROLLBACK || + it == TransactionState.TRANSACTION_STATE_COMMIT + ) { + error("Cannot join transaction cause, transaction \"$transactionId\" already ${it.name}") + } + transactionId + } .joinTransaction(undo) - .subscribeTransaction() - .watchTransaction() .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } @@ -56,18 +57,6 @@ abstract class AbstractTransactionManager( } } - private fun Mono.subscribeTransaction(): Mono { - return this.doOnSuccess { - transactionListener.subscribeStream(it) - .subscribeOn(Schedulers.parallel()) - .subscribe() - } - } - - private fun Mono.watchTransaction(): Mono { - return this.flatMap { transactionRetrySupporter.watchTransaction(it) } - } - final override fun rollback(transactionId: String, cause: String): Mono { return exists(transactionId) .publishTransaction(transaction { @@ -93,17 +82,13 @@ abstract class AbstractTransactionManager( final override fun exists(transactionId: String): Mono { return findAnyTransaction(transactionId) - .switchIfEmpty( - Mono.error { - IllegalStateException("Cannot find exists transaction id \"$transactionId\"") - } - ).mapTransactionId() + .mapTransactionId() .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } - protected abstract fun findAnyTransaction(transactionId: String): Mono + protected abstract fun findAnyTransaction(transactionId: String): Mono - protected fun Mono<*>.mapTransactionId(): Mono { + private fun Mono<*>.mapTransactionId(): Mono { return this.flatMap { Mono.deferContextual { Mono.just(it["transactionId"]) } } diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionRetrySupporter.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionRetrySupporter.kt index 7ad62a6..7f56262 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionRetrySupporter.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionRetrySupporter.kt @@ -21,7 +21,5 @@ abstract class AbstractTransactionRetrySupporter( .subscribe() } - abstract fun watchTransaction(transactionId: String): Mono - protected abstract fun handleOrphanTransaction(): Flux> } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index 69a5948..a816fe0 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -21,7 +21,6 @@ import kotlin.reflect.full.declaredMemberFunctions class RedisStreamTransactionDispatcher( private val applicationContext: ApplicationContext, private val reactiveRedisTemplate: ReactiveRedisTemplate, - private val redisStreamTransactionRemover: RedisStreamTransactionRemover, private val nodeGroup: String, ) : AbstractTransactionDispatcher() { @@ -67,7 +66,7 @@ class RedisStreamTransactionDispatcher( override fun findOwnTransaction(transaction: Transaction): Mono { return reactiveRedisTemplate.opsForStream() - .read(StreamOffset.create(transaction.id, ReadOffset.from("0"))) + .read(StreamOffset.create(STREAM_KEY, ReadOffset.from("0"))) .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } .filter { it.group == nodeGroup } .filter { hasUndo(it) } @@ -80,7 +79,7 @@ class RedisStreamTransactionDispatcher( override fun ack(transaction: Transaction, messageId: String): Mono> { return reactiveRedisTemplate.opsForStream() - .acknowledge(transaction.id, nodeGroup, messageId) + .acknowledge(STREAM_KEY, nodeGroup, messageId) .map { transaction to messageId } .switchIfEmpty( Mono.error { @@ -89,12 +88,9 @@ class RedisStreamTransactionDispatcher( ) } - override fun deleteElastic( - transaction: Transaction, - messageId: String - ) = redisStreamTransactionRemover.deleteElastic(transaction) - private companion object { + private const val STREAM_KEY = "NETX_STREAM" + private val notMatchedTransactionHandlerException = IllegalStateException("Cannot find matched Transaction handler") } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt index b388eca..5615cca 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionListener.kt @@ -14,7 +14,6 @@ import reactor.core.publisher.Flux import reactor.core.publisher.Mono import reactor.core.scheduler.Schedulers import kotlin.time.Duration.Companion.hours -import kotlin.time.Duration.Companion.milliseconds import kotlin.time.toJavaDuration class RedisStreamTransactionListener( @@ -31,26 +30,30 @@ class RedisStreamTransactionListener( private val receiver = StreamReceiver.create(connectionFactory, options) - override fun receive(transactionId: String): Flux> { - return createGroupIfNotExists(transactionId) + override fun receive(): Flux> { + return createGroupIfNotExists() .flatMap { receiver.receive( Consumer.from(nodeGroup, nodeName), - StreamOffset.create(transactionId, ReadOffset.from(">")) + StreamOffset.create(STREAM_KEY, ReadOffset.from(">")) ).publishOn(Schedulers.parallel()) .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) to it.id.value } } } - private fun createGroupIfNotExists(transactionId: String): Flux { + private fun createGroupIfNotExists(): Flux { return reactiveRedisTemplate.opsForStream() - .createGroup(transactionId, ReadOffset.from("0"), nodeGroup) + .createGroup(STREAM_KEY, ReadOffset.from("0"), nodeGroup) .onErrorResume { if (it.cause is RedisBusyException) { - return@onErrorResume Mono.just(transactionId) + return@onErrorResume Mono.just("OK") } throw it } .flatMapMany { Flux.just(it) } } + + private companion object { + private const val STREAM_KEY = "NETX_STREAM" + } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 531f5ee..96b91f8 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -1,46 +1,65 @@ package org.rooftop.netx.redis -import org.rooftop.netx.engine.AbstractTransactionListener +import org.redisson.api.RedissonReactiveClient import org.rooftop.netx.engine.AbstractTransactionManager import org.rooftop.netx.engine.AbstractTransactionRetrySupporter import org.rooftop.netx.idl.Transaction -import org.springframework.data.domain.Range +import org.rooftop.netx.idl.TransactionState import org.springframework.data.redis.connection.stream.Record import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono +import reactor.core.scheduler.Schedulers +import java.util.concurrent.TimeUnit class RedisStreamTransactionManager( nodeId: Int, nodeName: String, nodeGroup: String, - transactionListener: AbstractTransactionListener, transactionRetrySupporter: AbstractTransactionRetrySupporter, private val reactiveRedisTemplate: ReactiveRedisTemplate, + private val redissonReactiveClient: RedissonReactiveClient, ) : AbstractTransactionManager( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionListener = transactionListener, - transactionRetrySupporter = transactionRetrySupporter, ) { - override fun findAnyTransaction(transactionId: String): Mono { - return reactiveRedisTemplate.opsForStream() - .range(transactionId, Range.open("-", "+")) - .map { Transaction.parseFrom(it.value[DATA].toString().toByteArray()) } - .next() + override fun findAnyTransaction(transactionId: String): Mono { + return reactiveRedisTemplate + .opsForValue()[transactionId] + .switchIfEmpty( + Mono.error { + error("Cannot find exists transaction id \"$transactionId\"") + } + ) + .map { TransactionState.valueOf(String(it)) } } override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { return reactiveRedisTemplate.opsForStream() .add( Record.of(mapOf(DATA to transaction.toByteArray())) - .withStreamKey(transactionId) + .withStreamKey(STREAM_KEY) ) - .mapTransactionId() + .flatMap { + redissonReactiveClient.getLock("$transactionId-key") + .tryLock(10, TimeUnit.MINUTES) + } + .flatMap { + reactiveRedisTemplate.opsForValue() + .set(transactionId, transaction.state.name.toByteArray()) + } + .doFinally { + redissonReactiveClient.getLock("$transactionId-key") + .forceUnlock() + .subscribeOn(Schedulers.parallel()) + .subscribe() + } + .map { transactionId } } private companion object { private const val DATA = "data" + private const val STREAM_KEY = "NETX_STREAM" } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemover.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemover.kt deleted file mode 100644 index b7cb544..0000000 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemover.kt +++ /dev/null @@ -1,51 +0,0 @@ -package org.rooftop.netx.redis - -import org.rooftop.netx.idl.Transaction -import org.rooftop.netx.idl.TransactionState -import org.springframework.data.domain.Range -import org.springframework.data.redis.core.ReactiveRedisTemplate -import reactor.core.publisher.Mono -import reactor.core.scheduler.Schedulers -import reactor.util.retry.RetrySpec -import java.time.Duration - -class RedisStreamTransactionRemover( - private val nodeGroup: String, - private val reactiveRedisTemplate: ReactiveRedisTemplate, -) { - - fun deleteElastic(transaction: Transaction) { - Mono.just(transaction) - .filter { isTransactionEndState(transaction) } - .flatMap { - reactiveRedisTemplate.opsForStream() - .pending(transaction.id, nodeGroup, Range.closed("-", "+"), Long.MAX_VALUE) - .filter { - when (it.get().toList().isEmpty()) { - true -> true - false -> error(TRANSACTION_IS_PENDING_STATUS) - } - } - .retryWhen(retryIfTransactionPending) - .flatMap { - reactiveRedisTemplate.opsForSet() - .remove(nodeGroup, transaction.id.toByteArray()) - } - }.subscribeOn(Schedulers.parallel()) - .subscribe() - } - - private fun isTransactionEndState(transaction: Transaction): Boolean { - return transaction.state == TransactionState.TRANSACTION_STATE_ROLLBACK - || transaction.state == TransactionState.TRANSACTION_STATE_COMMIT - } - - companion object { - private const val TRANSACTION_IS_PENDING_STATUS = - "Transaction message remains in pending status." - private val retryIfTransactionPending = - RetrySpec.fixedDelay(Long.MAX_VALUE, Duration.ofMillis(3000)) - .jitter(1.0) - .filter { it.message == TRANSACTION_IS_PENDING_STATUS } - } -} diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt index 8497c99..fb2939c 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt @@ -34,9 +34,9 @@ class RedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionListener = redisStreamTransactionListener(), transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), + redissonReactiveClient = redissonReactiveClient(), ) @Bean @@ -48,7 +48,7 @@ class RedisTransactionConfigurer( nodeGroup = nodeGroup, nodeName = nodeName, reactiveRedisTemplate = reactiveRedisTemplate() - ) + ).also { it.subscribeStream() } @Bean @ConditionalOnProperty(prefix = "netx", name = ["mode"], havingValue = "redis") @@ -69,16 +69,7 @@ class RedisTransactionConfigurer( RedisStreamTransactionDispatcher( applicationContext = applicationContext, reactiveRedisTemplate = reactiveRedisTemplate(), - redisStreamTransactionRemover = redisStreamTransactionRemover(), - nodeGroup = nodeGroup, - ) - - @Bean - @ConditionalOnProperty(prefix = "netx", name = ["mode"], havingValue = "redis") - fun redisStreamTransactionRemover(): RedisStreamTransactionRemover = - RedisStreamTransactionRemover( nodeGroup = nodeGroup, - reactiveRedisTemplate = reactiveRedisTemplate(), ) @Bean @@ -98,11 +89,10 @@ class RedisTransactionConfigurer( fun redissonReactiveClient(): RedissonReactiveClient { val port: String = System.getProperty("netx.port") ?: port - return Redisson.create(Config() - .also { - it.useSingleServer() - .setAddress("redis://$host:$port") - }).reactive() + return Redisson.create(Config().also { + it.useSingleServer() + .setAddress("redis://$host:$port") + }).reactive() } @Bean diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionRetrySupporter.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionRetrySupporter.kt index 99a66f5..9627cb2 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionRetrySupporter.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionRetrySupporter.kt @@ -8,7 +8,6 @@ import org.springframework.data.domain.Range import org.springframework.data.redis.connection.RedisStreamCommands.XClaimOptions import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Flux -import reactor.core.publisher.Mono import reactor.core.scheduler.Schedulers import java.util.concurrent.TimeUnit @@ -23,16 +22,8 @@ class RedisTransactionRetrySupporter( private val lockKey: String = "$nodeGroup-key", ) : AbstractTransactionRetrySupporter(recoveryMilli) { - override fun watchTransaction(transactionId: String): Mono { - return reactiveRedisTemplate.opsForSet() - .add(nodeGroup, transactionId.toByteArray()) - .map { transactionId } - } - override fun handleOrphanTransaction(): Flux> { - return reactiveRedisTemplate.opsForSet() - .members(nodeGroup) - .flatMap { claimTransactions(String(it)) } + return claimTransactions() .publishOn(Schedulers.parallel()) .flatMap { (transaction, messageId) -> transactionDispatcher.dispatch(transaction, messageId) @@ -40,9 +31,9 @@ class RedisTransactionRetrySupporter( } } - private fun claimTransactions(transactionId: String): Flux> { + private fun claimTransactions(): Flux> { return reactiveRedisTemplate.opsForStream() - .pending(transactionId, nodeGroup, Range.closed("-", "+"), Long.MAX_VALUE) + .pending(STREAM_KEY, nodeGroup, Range.closed("-", "+"), Long.MAX_VALUE) .filter { it.get().toList().isNotEmpty() } .flatMap { pendingMessage -> redissonReactiveClient.getLock(lockKey) @@ -52,8 +43,10 @@ class RedisTransactionRetrySupporter( .flatMapMany { reactiveRedisTemplate.opsForStream() .claim( - transactionId, nodeGroup, nodeName, XClaimOptions - .minIdleMs(orphanMilli) + STREAM_KEY, + nodeGroup, + nodeName, + XClaimOptions.minIdleMs(orphanMilli) .ids(it.get().map { eachMessage -> eachMessage.id.value }.toList()) ) } @@ -70,4 +63,8 @@ class RedisTransactionRetrySupporter( .subscribe() } } + + private companion object { + private const val STREAM_KEY = "NETX_STREAM" + } } diff --git a/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt b/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt index 30501e1..c517bd9 100644 --- a/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt +++ b/src/test/kotlin/org/rooftop/netx/client/NetxTest.kt @@ -33,9 +33,6 @@ internal class NetxTest( row(10, 10), row(100, 100), row(1_000, 1_000), - row(10_000, 10_000), - row(100_000, 100_000), - row(1_000_000, 1_000_000), ) { commitLoadCount, rollbackLoadCount -> transactionReceiveStorage.clear() @@ -51,14 +48,12 @@ internal class NetxTest( netxClient.rollbackTransaction(transactionId, "") } - eventually(10.minutes) { + eventually(30.minutes) { transactionReceiveStorage.startCountShouldBe(commitLoadCount + rollbackLoadCount) transactionReceiveStorage.joinCountShouldBe(commitLoadCount + rollbackLoadCount) transactionReceiveStorage.commitCountShouldBe(commitLoadCount) transactionReceiveStorage.rollbackCountShouldBe(rollbackLoadCount) } - - Thread.sleep(Long.MAX_VALUE) } } }) diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt index e760edb..e7506b6 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt @@ -77,12 +77,6 @@ class NoAckRedisStreamTransactionDispatcher( override fun ack(transaction: Transaction, messageId: String): Mono> = Mono.just(transaction to messageId) - override fun deleteElastic( - transaction: Transaction, - messageId: String - ) { - } - private companion object { private val notMatchedTransactionHandlerException = IllegalStateException("Cannot find matched Transaction handler") diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt index 696ba6d..429eb40 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt @@ -34,9 +34,9 @@ class NoAckRedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionListener = redisStreamTransactionListener(), transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), + redissonReactiveClient = redissonReactiveClient(), ) @Bean @@ -48,7 +48,7 @@ class NoAckRedisTransactionConfigurer( nodeGroup = nodeGroup, nodeName = nodeName, reactiveRedisTemplate = reactiveRedisTemplate() - ) + ).also { it.subscribeStream() } @Bean @ConditionalOnProperty(prefix = "netx", name = ["mode"], havingValue = "redis") @@ -78,16 +78,7 @@ class NoAckRedisTransactionConfigurer( RedisStreamTransactionDispatcher( applicationContext = applicationContext, reactiveRedisTemplate = reactiveRedisTemplate(), - redisStreamTransactionRemover = redisStreamTransactionRemover(), - nodeGroup = nodeGroup, - ) - - @Bean - @ConditionalOnProperty(prefix = "netx", name = ["mode"], havingValue = "redis") - fun redisStreamTransactionRemover(): RedisStreamTransactionRemover = - RedisStreamTransactionRemover( nodeGroup = nodeGroup, - reactiveRedisTemplate = reactiveRedisTemplate(), ) @Bean @@ -108,11 +99,10 @@ class NoAckRedisTransactionConfigurer( val port: String = System.getProperty("netx.port") ?: port return Redisson.create( - Config() - .also { - it.useSingleServer() - .setAddress("redis://$host:$port") - }).reactive() + Config().also { + it.useSingleServer() + .setAddress("redis://$host:$port") + }).reactive() } @Bean diff --git a/src/test/kotlin/org/rooftop/netx/redis/RedisAssertions.kt b/src/test/kotlin/org/rooftop/netx/redis/RedisAssertions.kt index de5c675..2201887 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/RedisAssertions.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/RedisAssertions.kt @@ -14,7 +14,7 @@ internal class RedisAssertions( fun pendingMessageCountShouldBe(transactionId: String, count: Long) { val pendingMessageCount = reactiveRedisOperations.opsForStream() - .pending(transactionId, nodeGroup, Range.closed("-", "+"), Long.MAX_VALUE) + .pending("NETX_STREAM", nodeGroup, Range.closed("-", "+"), Long.MAX_VALUE) .map { it.get().toList().size } .block() diff --git a/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt b/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt index 0d9892c..816e554 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt @@ -3,7 +3,6 @@ package org.rooftop.netx.redis import io.kotest.assertions.nondeterministic.eventually import io.kotest.core.annotation.DisplayName import io.kotest.core.spec.style.DescribeSpec -import io.kotest.matchers.shouldBe import org.rooftop.netx.api.* import org.rooftop.netx.meta.EnableDistributedTransaction import org.springframework.test.context.ContextConfiguration diff --git a/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemoverTest.kt b/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemoverTest.kt deleted file mode 100644 index de4f86c..0000000 --- a/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionRemoverTest.kt +++ /dev/null @@ -1,66 +0,0 @@ -package org.rooftop.netx.redis - -import io.kotest.assertions.nondeterministic.eventually -import io.kotest.core.annotation.DisplayName -import io.kotest.core.spec.style.DescribeSpec -import org.rooftop.netx.api.TransactionManager -import org.rooftop.netx.meta.EnableDistributedTransaction -import org.springframework.test.context.ContextConfiguration -import org.springframework.test.context.TestPropertySource -import reactor.core.scheduler.Schedulers -import kotlin.time.Duration.Companion.minutes -import kotlin.time.Duration.Companion.seconds - -@EnableDistributedTransaction -@ContextConfiguration( - classes = [ - RedisContainer::class, - RedisAssertions::class, - TransactionHandlerAssertions::class, - ] -) -@TestPropertySource("classpath:application.properties") -@DisplayName("RedisStreamTransactionRemover 클래스의") -internal class RedisStreamTransactionRemoverTest( - private val redisAssertions: RedisAssertions, - private val transactionManager: TransactionManager, - private val transactionHandlerAssertions: TransactionHandlerAssertions, -) : DescribeSpec({ - - beforeEach { - transactionHandlerAssertions.clear() - } - - describe("handleTransactionCommitEvent 메소드는") { - context("TransactionCommitEvent 가 발행되면,") { - val transactionId = transactionManager.start("RedisStreamTransactionRemoverTest") - .block()!! - - it("Transaction 을 retry watch 대기열에서 삭제한다.") { - transactionManager.commit(transactionId) - .subscribeOn(Schedulers.parallel()) - .subscribe() - - eventually(5.minutes) { - transactionHandlerAssertions.commitCountShouldBe(1) - redisAssertions.retryTransactionShouldBeNotExists(transactionId) - } - } - } - - context("TransactionRollbackEvent 가 발행되면,") { - val transactionId = transactionManager.start("RedisStreamTransactionRemoverTest") - .block()!! - - it("Transaction 을 retry watch 대기열에서 삭제한다.") { - transactionManager.rollback(transactionId, "rollback occured for test").block() - - eventually(5.minutes) { - transactionHandlerAssertions.rollbackCountShouldBe(1) - redisAssertions.retryTransactionShouldBeNotExists(transactionId) - } - } - } - } -} -)