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
51 changes: 31 additions & 20 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,15 @@
> Distributed transaction library based on Choreography

<br>

![version 0.1.1](https://img.shields.io/badge/version-0.1.1-black?labelColor=black&style=flat-square) ![jdk 17](https://img.shields.io/badge/jdk-17-orange?labelColor=black&style=flat-square)
![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"/>

Choreography 방식으로 구현된 분산 트랜잭션 라이브러리 입니다.
`Netx` 는 다음 기능을 제공합니다.
`Netx` 는 다음 기능을 제공합니다.

1. [Reactor](https://projectreactor.io/) 기반의 완전한 비동기 트랜잭션 관리
2. Redis-stream 기반의 트랜잭션 관리
2. Redis-stream 기반의 트랜잭션 관리
3. 여러 노드가 중복 트랜잭션 이벤트를 수신하는 문제 방지
4. `At Least Once` 방식의 메시지 전달 보장

Expand All @@ -21,7 +21,7 @@ Netx는 스프링 환경에서 사용할 수 있으며, 아래와 같이 `@AutoC

```kotlin
@SpringBootApplication
@AutoConfigureRedisTransaction
@AutoConfigureDistributedTransaction
@EnableAutoConfiguration(exclude = [RedisReactiveAutoConfiguration::class])
class Application {

Expand All @@ -34,17 +34,21 @@ class Application {
}
```

`@AutoconfigureRedisTransaction` 어노테이션으로 자동 구성할 경우 netx는 아래 프로퍼티를 사용해 메시지 큐와 커넥션을 맺습니다.
`@AutoConfigureDistributedTransaction` 어노테이션으로 자동 구성할 경우 netx는 아래 프로퍼티를 사용해 메시지 큐와 커넥션을 맺습니다.

#### Properties
| key | example | description |
|------------------|----------|------------------------------------------------------------------------------------------------------------------------------------|
| **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`_ 에 참여할 서버의 이름입니다. 같은 그룹내에 중복된 이름이 존재하면 안됩니다. |

| 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 입니다. |

### Usage example

Expand All @@ -55,21 +59,27 @@ fun pay(param: Any): Mono<Any> {
return transactionManager.start("paid=1000") // Start distributed transaction and publish transaction start event
.flatMap { transactionId ->
service.pay(param)
.doOnError { throwable ->
transactionManager.rollback(transactionId, throwable.message) // Publish rollback event to all transaction joined node
.doOnError { throwable ->
transactionManager.rollback(
transactionId,
throwable.message
) // Publish rollback event to all transaction joined node
}
}.doOnSuccess { transactionId ->
transactionManager.commit(transactionId) // Publish commit event to all transaction joined node
}
}
```

#### Scenario2. Join order transaction

```kotlin
fun order(param: Any): Mono<Any> {
return transactionManager.join(param.transactionId, "orderId=1:state=PENDING") // join exists distributed transaction and publish transaction join event
.flatMap { transactionId ->
return transactionManager.join(
param.transactionId,
"orderId=1:state=PENDING"
) // join exists distributed transaction and publish transaction join event
.flatMap { transactionId ->
service.order(param)
.doOnError { throwable ->
transactionManager.rollback(transactionId, throwable.message)
Expand All @@ -79,7 +89,7 @@ fun order(param: Any): Mono<Any> {
}
}
```

#### Scenario3. Check exists transaction

```kotlin
Expand All @@ -90,7 +100,8 @@ fun exists(param: Any): Mono<Any> {

#### Scenario4. Handle transaction event

다른 분산서버가 (혹은 자기자신이) transactionManager를 통해서 트랜잭션을 시작하거나 트랜잭션 상태를 변경했을때, 호출한 메소드에 맞는 트랜잭션 이벤트를 발행합니다.
다른 분산서버가 (혹은 자기자신이) transactionManager를 통해서 트랜잭션을 시작하거나 트랜잭션 상태를 변경했을때, 호출한 메소드에 맞는 트랜잭션 이벤트를
발행합니다.
이 이벤트들을 핸들링 함으로써, 다른서버에서 발생한 에러등을 수신하고 롤백할 수 있습니다.

```kotlin
Expand Down
5 changes: 4 additions & 1 deletion gradle.properties
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ kotlin.code.style=official

### Project ###
group=org.rooftop.netx
version=0.1.0
version=0.1.2
compatibility=17

### Protobuf ###
Expand Down Expand Up @@ -36,3 +36,6 @@ lettuceVersion=6.3.0.RELEASE

### TestContainer ###
testContainerVersion=1.19.3

### Jackson ###
jacksonVersion=2.16.1
2 changes: 2 additions & 0 deletions gradle/core.gradle
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
dependencies {
implementation "com.github.f4b6a3:tsid-creator:${snowflakeVersion}"
implementation "com.fasterxml.jackson.core:jackson-databind:${jacksonVersion}"
implementation "com.fasterxml.jackson.module:jackson-module-parameter-names:${jacksonVersion}"
}
2 changes: 1 addition & 1 deletion idl
6 changes: 3 additions & 3 deletions src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,11 @@ import reactor.core.publisher.Mono

interface TransactionManager {

fun start(replay: String): Mono<String>
fun start(undo: String): Mono<String>

fun exists(transactionId: String): Mono<String>
fun join(transactionId: String, undo: String): Mono<String>

fun join(transactionId: String, replay: String): Mono<String>
fun exists(transactionId: String): Mono<String>

fun commit(transactionId: String): Mono<String>

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ package org.rooftop.netx.api

data class TransactionRollbackEvent(
val transactionId: String,
val replay: String,
val nodeName: String,
val cause: String?,
val undoState: String,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
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)
annotation class AutoConfigureDistributedTransaction
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,14 @@ import reactor.core.publisher.Flux
import reactor.core.publisher.Mono

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

@EventListener(SubscribeTransactionEvent::class)
fun subscribeStream(event: SubscribeTransactionEvent): Flux<Transaction> {
return receive(event).dispatch()
return receive(event)
.dispatch()
}

protected abstract fun receive(event: SubscribeTransactionEvent): Flux<Transaction>
Expand Down Expand Up @@ -50,18 +52,20 @@ abstract class AbstractTransactionDispatcher(
.doOnNext { eventPublisher.publish(TransactionCommitEvent(it.id, it.serverId)) }
}

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

private fun publishStart(it: Transaction): Mono<Transaction> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,39 +11,40 @@ abstract class AbstractTransactionManager(
private val nodeName: String,
private val eventPublisher: EventPublisher,
private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId),
private val undoManager: UndoManager,
) : TransactionManager {

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

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

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

private fun Mono<String>.joinTransaction(replay: String): Mono<String> {
private fun Mono<String>.joinTransaction(): Mono<String> {
return flatMap { transactionId ->
publishTransaction(transactionId, transaction {
id = transactionId
serverId = nodeName
this.replay = replay
state = TransactionState.TRANSACTION_STATE_JOIN
})
}
Expand All @@ -55,6 +56,10 @@ abstract class AbstractTransactionManager(
}
}

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

final override fun rollback(transactionId: String, cause: String): Mono<String> {
return exists(transactionId)
.publishTransaction(transaction {
Expand Down
13 changes: 13 additions & 0 deletions src/main/kotlin/org/rooftop/netx/engine/UndoManager.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
package org.rooftop.netx.engine

import reactor.core.publisher.Mono

interface UndoManager {

fun save(transactionId: String, undo: String): Mono<String>

fun find(transactionId: String): Mono<String>

fun delete(transactionId: String): Mono<Boolean>

}

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package org.rooftop.netx.redis

import org.rooftop.netx.engine.AbstractTransactionDispatcher
import org.rooftop.netx.engine.SubscribeTransactionEvent
import org.rooftop.netx.engine.UndoManager
import org.rooftop.netx.idl.Transaction
import org.springframework.context.ApplicationEventPublisher
import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory
Expand All @@ -18,10 +19,11 @@ import kotlin.time.toJavaDuration
class RedisStreamTransactionDispatcher(
eventPublisher: ApplicationEventPublisher,
connectionFactory: ReactiveRedisConnectionFactory,
undoManager: UndoManager,
private val streamGroup: String,
private val nodeName: String,
private val reactiveRedisTemplate: ReactiveRedisTemplate<String, ByteArray>,
) : AbstractTransactionDispatcher(SpringEventPublisher(eventPublisher)) {
) : AbstractTransactionDispatcher(undoManager, SpringEventPublisher(eventPublisher)) {

private val options = StreamReceiver.StreamReceiverOptions.builder()
.pollTimeout(1.hours.toJavaDuration())
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package org.rooftop.netx.redis

import org.rooftop.netx.engine.AbstractTransactionManager
import org.rooftop.netx.engine.UndoManager
import org.rooftop.netx.idl.Transaction
import org.springframework.context.ApplicationEventPublisher
import org.springframework.data.domain.Range
Expand All @@ -12,8 +13,14 @@ class RedisStreamTransactionManager(
nodeId: Int,
nodeName: String,
applicationEventPublisher: ApplicationEventPublisher,
undoManager: UndoManager,
private val reactiveRedisTemplate: ReactiveRedisTemplate<String, ByteArray>,
) : AbstractTransactionManager(nodeId, nodeName, SpringEventPublisher(applicationEventPublisher)) {
) : AbstractTransactionManager(
nodeId,
nodeName,
SpringEventPublisher(applicationEventPublisher),
undoManager = undoManager
) {

override fun findAnyTransaction(transactionId: String): Mono<Transaction> {
return reactiveRedisTemplate.opsForStream<String, ByteArray>()
Expand Down
Loading