From e818366fd88962673cb0b53bf105c33f5b405777 Mon Sep 17 00:00:00 2001 From: devxb Date: Mon, 19 Feb 2024 00:31:51 +0900 Subject: [PATCH 1/2] =?UTF-8?q?=EC=9E=90=EC=8B=A0=EC=9D=B4=20=EB=B0=9C?= =?UTF-8?q?=ED=96=89=ED=95=9C=20transaction=EC=9D=84=20=EC=B0=BE=EC=A7=80?= =?UTF-8?q?=20=EB=AA=BB=ED=95=98=EB=8A=94=20=EB=B2=84=EA=B7=B8=EB=A5=BC=20?= =?UTF-8?q?=EC=88=98=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 +-- .../redis/RedisStreamTransactionDispatcher.kt | 20 +++------ .../redis/RedisStreamTransactionManager.kt | 44 ++++++++++--------- .../netx/redis/RedisTransactionConfigurer.kt | 2 - .../NoAckRedisStreamTransactionDispatcher.kt | 14 +++--- .../redis/NoAckRedisTransactionConfigurer.kt | 2 - 6 files changed, 40 insertions(+), 48 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt index 3a351e6..b4f5efc 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt @@ -61,14 +61,14 @@ abstract class AbstractTransactionDispatcher { ) ) - TransactionState.TRANSACTION_STATE_ROLLBACK -> findOwnTransaction(transaction) + TransactionState.TRANSACTION_STATE_ROLLBACK -> findOwnUndo(transaction) .map { TransactionRollbackEvent( transaction.id, transaction.serverId, transaction.group, transaction.cause, - it.undo, + it, ) } @@ -76,7 +76,7 @@ abstract class AbstractTransactionDispatcher { } } - protected abstract fun findOwnTransaction(transaction: Transaction): Mono + protected abstract fun findOwnUndo(transaction: Transaction): Mono protected abstract fun ack( transaction: Transaction, diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index a816fe0..7263cd0 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -10,8 +10,6 @@ import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState import org.rooftop.netx.meta.TransactionHandler import org.springframework.context.ApplicationContext -import org.springframework.data.redis.connection.stream.ReadOffset -import org.springframework.data.redis.connection.stream.StreamOffset import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono import kotlin.reflect.KClass @@ -64,19 +62,15 @@ class RedisStreamTransactionDispatcher( } } - override fun findOwnTransaction(transaction: Transaction): Mono { - return reactiveRedisTemplate.opsForStream() - .read(StreamOffset.create(STREAM_KEY, ReadOffset.from("0"))) - .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } - .filter { it.group == nodeGroup } - .filter { hasUndo(it) } - .next() + override fun findOwnUndo(transaction: Transaction): Mono { + return reactiveRedisTemplate.opsForHash()[transaction.id, nodeGroup] + .switchIfEmpty( + Mono.error { + error("Cannot find undo state in transaction hashes key \"${transaction.id}\"") + } + ) } - private fun hasUndo(transaction: Transaction): Boolean = - transaction.state == TransactionState.TRANSACTION_STATE_JOIN - || transaction.state == TransactionState.TRANSACTION_STATE_START - override fun ack(transaction: Transaction, messageId: String): Mono> { return reactiveRedisTemplate.opsForStream() .acknowledge(STREAM_KEY, nodeGroup, messageId) diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 96b91f8..429071b 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -1,23 +1,17 @@ package org.rooftop.netx.redis -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.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, - transactionRetrySupporter: AbstractTransactionRetrySupporter, + private val nodeGroup: String, private val reactiveRedisTemplate: ReactiveRedisTemplate, - private val redissonReactiveClient: RedissonReactiveClient, ) : AbstractTransactionManager( nodeId = nodeId, nodeName = nodeName, @@ -26,13 +20,13 @@ class RedisStreamTransactionManager( override fun findAnyTransaction(transactionId: String): Mono { return reactiveRedisTemplate - .opsForValue()[transactionId] + .opsForHash()[transactionId, STATE_KEY] .switchIfEmpty( Mono.error { error("Cannot find exists transaction id \"$transactionId\"") } ) - .map { TransactionState.valueOf(String(it)) } + .map { TransactionState.valueOf(it) } } override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { @@ -42,24 +36,32 @@ class RedisStreamTransactionManager( .withStreamKey(STREAM_KEY) ) .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() + if (hasUndo(transaction)) { + return@flatMap reactiveRedisTemplate.opsForHash() + .putAll( + transactionId, mapOf( + STATE_KEY to transaction.state.name, + nodeGroup to transaction.undo + ) + ) + } + reactiveRedisTemplate.opsForHash() + .putAll( + transactionId, mapOf( + STATE_KEY to transaction.state.name, + ) + ) } .map { transactionId } } + private fun hasUndo(transaction: Transaction): Boolean = + transaction.state == TransactionState.TRANSACTION_STATE_JOIN + || transaction.state == TransactionState.TRANSACTION_STATE_START + private companion object { private const val DATA = "data" private const val STREAM_KEY = "NETX_STREAM" + private const val STATE_KEY = "TX_STATE" } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt index fb2939c..7d4fe88 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt @@ -34,9 +34,7 @@ class RedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), - redissonReactiveClient = redissonReactiveClient(), ) @Bean diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt index e7506b6..c93b783 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt @@ -61,13 +61,13 @@ class NoAckRedisStreamTransactionDispatcher( } } - override fun findOwnTransaction(transaction: Transaction): Mono { - return reactiveRedisTemplate.opsForStream() - .read(StreamOffset.create(transaction.id, ReadOffset.from("0"))) - .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } - .filter { it.group == nodeGroup } - .filter { hasUndo(it) } - .next() + override fun findOwnUndo(transaction: Transaction): Mono { + return reactiveRedisTemplate.opsForHash()[transaction.id, nodeGroup] + .switchIfEmpty( + Mono.error { + error("Cannot find undo state in transaction hashes key \"${transaction.id}\"") + } + ) } private fun hasUndo(transaction: Transaction): Boolean = diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt index 429eb40..88cd6c9 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt @@ -34,9 +34,7 @@ class NoAckRedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), - redissonReactiveClient = redissonReactiveClient(), ) @Bean From 2fad59e73679380456cf8f158e9c1285d711644f Mon Sep 17 00:00:00 2001 From: devxb Date: Mon, 19 Feb 2024 00:33:34 +0900 Subject: [PATCH 2/2] =?UTF-8?q?refactor:=20=EC=82=AC=EC=9A=A9=ED=95=98?= =?UTF-8?q?=EC=A7=80=20=EC=95=8A=EB=8A=94=20=EB=A9=94=EC=86=8C=EB=93=9C?= =?UTF-8?q?=EB=A5=BC=20=EC=82=AD=EC=A0=9C=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../netx/redis/NoAckRedisStreamTransactionDispatcher.kt | 4 ---- 1 file changed, 4 deletions(-) diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt index c93b783..00a060b 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt @@ -70,10 +70,6 @@ class NoAckRedisStreamTransactionDispatcher( ) } - private fun hasUndo(transaction: Transaction): Boolean = - transaction.state == TransactionState.TRANSACTION_STATE_JOIN - || transaction.state == TransactionState.TRANSACTION_STATE_START - override fun ack(transaction: Transaction, messageId: String): Mono> = Mono.just(transaction to messageId)