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
23 changes: 12 additions & 11 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@
> Distributed transaction library based on Choreography

<br>


![version 0.1.2](https://img.shields.io/badge/version-0.1.2-black?labelColor=black&style=flat-square) ![jdk 17](https://img.shields.io/badge/jdk-17-orange?labelColor=black&style=flat-square)

<img src = "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/rooftop-MSA/Netx/assets/62425964/5082ef20-10ad-4b6b-bff8-7e78a0f9e01f" width="500" align="right"/>
Expand Down Expand Up @@ -38,17 +40,16 @@ class Application {

#### Properties

| key | example | description |
|--------------------|-----------|------------------------------------------------------------------------------------------------------------------------------------|
| **netx.mode** | redis | 트랜잭션 관리에 사용할 메시지 큐 구현체의 mode 입니다. |
| **netx.host** | localhost | 트랜잭션 관리에 사용할 메시지 큐 의 host url 입니다. (ex. redis host) |
| **netx.port** | 6379 | 트랜잭션 관리에 사용할 메시지 큐의 port 입니다. |
| **netx.group** | pay-group | 분산 노드의 그룹입니다. 트랜잭션 이벤트는 같은 그룹내 하나의 노드로만 전송됩니다. |
| **netx.node-id** | 1 | id 생성에 사용될 식별자입니다. 모든 서버는 반드시 다른 id를 할당받아야 하며, 1~256 만큼의 id를 설정할 수 있습니다. _`중복된 id 생성을 방지하기위해 twitter snowflake 알고리즘으로 id를 생성합니다.`_ |
| **netx.node-name** | pay-1 | _`$netx.group`_ 에 참여할 서버의 이름입니다. 같은 그룹내에 중복된 이름이 존재하면 안됩니다. |
| **netx.undo.mode** | redis | 트랜잭션 undo 상태 저장에 사용할 저장소 구현체의 mode 입니다. |
| **netx.undo.host** | localhost | 트랜잭션 undo 상태 저장에 사용할 저장소의 host url 입니다. |
| **netx.undo.port** | 6380 | 트랜잭션 undo 상태 저장에 사용할 저장소의 port 입니다. |
| key | example | description |
|-------------------------|-----------|----------------------------------------------------------------------------------------------------------------------------------------------------|
| **netx.mode** | redis | 트랜잭션 관리에 사용할 메시지 큐 구현체의 mode 입니다. |
| **netx.host** | localhost | 트랜잭션 관리에 사용할 메시지 큐 의 host url 입니다. (ex. redis host) |
| **netx.port** | 6379 | 트랜잭션 관리에 사용할 메시지 큐의 port 입니다. |
| **netx.group** | pay-group | 분산 노드의 그룹입니다. 트랜잭션 이벤트는 같은 그룹내 하나의 노드로만 전송됩니다. |
| **netx.node-id** | 1 | id 생성에 사용될 식별자입니다. 모든 서버는 반드시 다른 id를 할당받아야 하며, 1~256 만큼의 id를 설정할 수 있습니다. _`중복된 id 생성을 방지하기위해 twitter snowflake 알고리즘으로 id를 생성합니다.`_ |
| **netx.node-name** | pay-1 | _`netx.group`_ 에 참여할 서버의 이름입니다. 같은 그룹내에 중복된 이름이 존재하면 안됩니다. |
| **netx.recovery-milli** | 60000 | _`netx.recovery-milli`_ 마다 _`netx.orphan-milli`_ 동안 처리 되지 않는 트랜잭션을 찾아 재실행합니다. 기본값은 60000(60초) 입니다. |
| **netx.orphan-milli** | 10000 | 트랜잭션이 PENDING 상태가 되었지만 orphan-milli가 지나도 ACK 상태가 되지 않는경우 다른 노드에게 처리를 위임합니다. 기본값은 10000(10초) 입니다. |

### Usage example

Expand Down
2 changes: 1 addition & 1 deletion build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ publishing {
}
}

apply from: "gradle/mq.gradle"
apply from: "gradle/db.gradle"
apply from: "gradle/test.gradle"
apply from: "gradle/core.gradle"
apply from: "gradle/sonar.gradle"
Expand Down
3 changes: 3 additions & 0 deletions gradle.properties
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,9 @@ snowflakeVersion=5.2.5
### Lettuce ###
lettuceVersion=6.3.0.RELEASE

### Redisson ###
redissonVersion=3.26.0

### TestContainer ###
testContainerVersion=1.19.3

Expand Down
1 change: 1 addition & 0 deletions gradle/mq.gradle → gradle/db.gradle
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
dependencies {
implementation "io.lettuce:lettuce-core:${lettuceVersion}"
implementation "org.redisson:redisson:${redissonVersion}"
}
2 changes: 1 addition & 1 deletion idl
Original file line number Diff line number Diff line change
Expand Up @@ -4,5 +4,5 @@ data class TransactionRollbackEvent(
val transactionId: String,
val nodeName: String,
val cause: String?,
val undoState: String,
val undo: String,
)
Original file line number Diff line number Diff line change
@@ -1,10 +1,9 @@
package org.rooftop.netx.autoconfig

import org.rooftop.netx.redis.RedisTransactionConfigurer
import org.rooftop.netx.redis.RedisUndoConfigurer
import org.springframework.boot.autoconfigure.ImportAutoConfiguration

@Target(AnnotationTarget.CLASS)
@Retention(AnnotationRetention.RUNTIME)
@ImportAutoConfiguration(RedisUndoConfigurer::class, RedisTransactionConfigurer::class)
@ImportAutoConfiguration(RedisTransactionConfigurer::class)
annotation class AutoConfigureDistributedTransaction
Original file line number Diff line number Diff line change
Expand Up @@ -6,39 +6,43 @@ import org.rooftop.netx.api.TransactionRollbackEvent
import org.rooftop.netx.api.TransactionStartEvent
import org.rooftop.netx.idl.Transaction
import org.rooftop.netx.idl.TransactionState
import org.springframework.context.event.EventListener
import org.springframework.context.ApplicationEventPublisher
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono

abstract class AbstractTransactionDispatcher(
private val undoManager: UndoManager,
private val eventPublisher: EventPublisher,
private val eventPublisher: ApplicationEventPublisher,
) {

@EventListener(SubscribeTransactionEvent::class)
fun subscribeStream(event: SubscribeTransactionEvent): Flux<Transaction> {
return receive(event)
.dispatch()
fun subscribeStream(transactionId: String): Flux<Pair<Transaction, String>> {
return receive(transactionId)
.flatMap { dispatchAndAck(it.first, it.second) }
}

protected abstract fun receive(event: SubscribeTransactionEvent): Flux<Transaction>
protected abstract fun receive(transactionId: String): Flux<Pair<Transaction, String>>

private fun Flux<Transaction>.dispatch(): Flux<Transaction> {
return this.flatMap {
when (it.state) {
TransactionState.TRANSACTION_STATE_JOIN -> publishJoin(it)
TransactionState.TRANSACTION_STATE_COMMIT -> publishCommit(it)
TransactionState.TRANSACTION_STATE_ROLLBACK -> publishRollback(it)
TransactionState.TRANSACTION_STATE_START -> publishStart(it)
else -> error("Cannot find matched transaction state \"${it.state}\"")
}
fun dispatchAndAck(transaction: Transaction, messageId: String): Flux<Pair<Transaction, String>> {
return Flux.just(transaction to messageId)
.dispatch()
.ack()
}

private fun Flux<Pair<Transaction, String>>.dispatch(): Flux<Pair<Transaction, String>> {
return this.flatMap { (transaction, messageId) ->
when (transaction.state) {
TransactionState.TRANSACTION_STATE_JOIN -> publishJoin(transaction)
TransactionState.TRANSACTION_STATE_COMMIT -> publishCommit(transaction)
TransactionState.TRANSACTION_STATE_ROLLBACK -> publishRollback(transaction)
TransactionState.TRANSACTION_STATE_START -> publishStart(transaction)
else -> error("Cannot find matched transaction state \"${transaction.state}\"")
}.map { transaction to messageId }
}
}

private fun publishJoin(it: Transaction): Mono<Transaction> {
return Mono.just(it)
.doOnNext {
eventPublisher.publish(
eventPublisher.publishEvent(
TransactionJoinEvent(
it.id,
it.serverId
Expand All @@ -49,29 +53,32 @@ abstract class AbstractTransactionDispatcher(

private fun publishCommit(it: Transaction): Mono<Transaction> {
return Mono.just(it)
.doOnNext { eventPublisher.publish(TransactionCommitEvent(it.id, it.serverId)) }
.doOnNext { eventPublisher.publishEvent(TransactionCommitEvent(it.id, it.serverId)) }
}

private fun publishRollback(transaction: Transaction): Mono<Transaction> {
return undoManager.find(transaction.id)
return findOwnTransaction(transaction)
.doOnNext {
eventPublisher.publish(
eventPublisher.publishEvent(
TransactionRollbackEvent(
transaction.id,
transaction.serverId,
transaction.cause,
it
it.undo
)
)
}
.flatMap { undoManager.delete(transaction.id) }
.map { transaction }
}

protected abstract fun findOwnTransaction(transaction: Transaction): Mono<Transaction>

private fun publishStart(it: Transaction): Mono<Transaction> {
return Mono.just(it)
.doOnNext {
eventPublisher.publish(TransactionStartEvent(it.id, it.serverId))
eventPublisher.publishEvent(TransactionStartEvent(it.id, it.serverId))
}
}

protected abstract fun Flux<Pair<Transaction, String>>.ack(): Flux<Pair<Transaction, String>>
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,66 +5,75 @@ 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 eventPublisher: EventPublisher,
private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId),
private val undoManager: UndoManager,
private val transactionDispatcher: AbstractTransactionDispatcher,
private val transactionRetrySupporter: AbstractTransactionRetrySupporter,
) : TransactionManager {

final override fun start(undo: String): Mono<String> {
return startTransaction()
return startTransaction(undo)
.subscribeTransaction()
.saveUndoState(undo)
.watchTransaction()
.contextWrite { it.put(CONTEXT_TX_KEY, transactionIdGenerator.generate()) }
}

private fun startTransaction(): Mono<String> {
private fun startTransaction(undo: String): Mono<String> {
return Mono.deferContextual<String> { Mono.just(it[CONTEXT_TX_KEY]) }
.flatMap { transactionId ->
publishTransaction(transactionId, transaction {
id = transactionId
serverId = nodeName
group = nodeGroup
this.state = TransactionState.TRANSACTION_STATE_START
this.undo = undo
})
}
}

final override fun join(transactionId: String, undo: String): Mono<String> {
return exists(transactionId)
.joinTransaction()
.joinTransaction(undo)
.subscribeTransaction()
.saveUndoState(undo)
.watchTransaction()
.contextWrite { it.put(CONTEXT_TX_KEY, transactionId) }
}

private fun Mono<String>.joinTransaction(): Mono<String> {
private fun Mono<String>.joinTransaction(undo: String): Mono<String> {
return flatMap { transactionId ->
publishTransaction(transactionId, transaction {
id = transactionId
serverId = nodeName
group = nodeGroup
state = TransactionState.TRANSACTION_STATE_JOIN
this.undo = undo
})
}
}

private fun Mono<String>.subscribeTransaction(): Mono<String> {
return this.doOnSuccess {
eventPublisher.publish(SubscribeTransactionEvent(it))
transactionDispatcher.subscribeStream(it)
.subscribeOn(Schedulers.parallel())
.subscribe()
}
}

private fun Mono<String>.saveUndoState(undo: String): Mono<String> {
return this.flatMap { undoManager.save(it, undo) }
private fun Mono<String>.watchTransaction(): Mono<String> {
return this.flatMap { transactionRetrySupporter.watchTransaction(it) }
}

final override fun rollback(transactionId: String, cause: String): Mono<String> {
return exists(transactionId)
.publishTransaction(transaction {
id = transactionId
serverId = nodeName
group = nodeGroup
state = TransactionState.TRANSACTION_STATE_ROLLBACK
this.cause = cause
})
Expand All @@ -76,6 +85,7 @@ abstract class AbstractTransactionManager(
.publishTransaction(transaction {
id = transactionId
serverId = nodeName
group = nodeGroup
state = TransactionState.TRANSACTION_STATE_COMMIT
})
.contextWrite { it.put(CONTEXT_TX_KEY, transactionId) }
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
package org.rooftop.netx.engine

import org.rooftop.netx.idl.Transaction
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.scheduler.Schedulers
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.toJavaDuration

abstract class AbstractTransactionRetrySupporter(
recoveryMilli: Long,
) {

init {
Flux.interval(recoveryMilli.milliseconds.toJavaDuration())
.publishOn(Schedulers.parallel())
.flatMap { handleOrphanTransaction() }
.subscribe()
}

abstract fun watchTransaction(transactionId: String): Mono<String>

protected abstract fun handleOrphanTransaction(): Flux<Pair<Transaction, String>>
}
6 changes: 0 additions & 6 deletions src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt

This file was deleted.

This file was deleted.

13 changes: 0 additions & 13 deletions src/main/kotlin/org/rooftop/netx/engine/UndoManager.kt

This file was deleted.

Loading