Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
4cc9929
feat: transaction api 를 정의한다
devxb Feb 2, 2024
5fac3ae
feat: 실제 동작이 구현되어있는 engine 레이어를 정의하고 동작을 구현한다
devxb Feb 2, 2024
e1de4fa
feat: Transaction의 상태별로 event 를 발행하게 한다
devxb Feb 2, 2024
fed0137
refactor: Transaction start 혹은 join일때, transaction을 구독하도록한다
devxb Feb 2, 2024
7d5226a
feat: redis-stream 기반의 transation-manager 를 구현한다
devxb Feb 2, 2024
ba67fe8
feat: AutoConfig 를 정의한다
devxb Feb 2, 2024
9b3dbf9
fix: TransactionIdGenerator에서 존재하지 않는 클래스 구현하고 있는 버그 수정
devxb Feb 2, 2024
b69b3a8
refactor: transaction id 를 node-id 기반으로 생성하도록 수정
devxb Feb 2, 2024
a6a5c05
refactor: 자동구성이 어노테이션 기반으로 동작 가능하도록 수정한다
devxb Feb 2, 2024
fb3ba4b
refactor: TransactionEvent들에 어떤 서버에서 발행되었는지 확인할 수 있도록 nodeName 필드를 추가한다
devxb Feb 2, 2024
fc4e076
refactor: TransactionDispatcher를 추상화시키고, 로직을 따르도록 강제화 한다
devxb Feb 2, 2024
8e176f9
fix: RedisStream 구독을 poll 방식이 아닌 subscirbe 방식으로 수정하고 처음부터 모든 데이터를 읽어오…
devxb Feb 2, 2024
84a70fc
refactor: replay 필드가 rollback 이벤트에만 존재하도록 수정한다
devxb Feb 4, 2024
c4e9dcb
refactor: exists 메소드 구현을 engine레이어로 올리고, context를 세팅한다
devxb Feb 4, 2024
6ffdf9f
refactor: transformTransactionId 이름을 mapTransactionId로 변경한다
devxb Feb 4, 2024
583c92d
test: RedisStreamTransactionManager의 Test를 작성한다
devxb Feb 4, 2024
968d766
build: jitpack 배포 플러그인을 설정한다
devxb Feb 4, 2024
84edd4f
docs: 사용법과 다운로드 방법을 작성한다
devxb Feb 4, 2024
a6baa7d
docs: 분산 트랜잭션 작동과정 gif를 추가한다
devxb Feb 4, 2024
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
123 changes: 123 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,126 @@
# Netx <img src="https://avatars.githubusercontent.com/u/149151221?s=200&v=4" height = 100 align = left>

> Distributed transaction library based on Choreography

![version 0.1.0](https://img.shields.io/badge/version-0.1.0-black?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` 는 다음 기능을 제공합니다.
1. [Reactor](https://projectreactor.io/) 기반의 완전한 비동기 트랜잭션 관리
2. Redis-stream 기반의 트랜잭션 관리
3. 여러 노드가 중복 트랜잭션 이벤트를 수신하는 문제 방지

## How to use

Netx는 스프링 환경에서 사용할 수 있으며, 아래와 같이 `@AutoConfigureRedisTransaction` 어노테이션을 붙이는것으로 손쉽게 사용할 수 있습니다.

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

companion object {
@JvmStatic
fun main(vararg args: String) {
SpringApplication.run(Application::class.java, *args)
}
}
}
```

`@AutoconfigureRedisTransaction` 어노테이션으로 자동 구성할 경우 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`_ 에 참여할 서버의 이름입니다. 같은 그룹내에 중복된 이름이 존재하면 안됩니다. |


### Usage example

#### Scenario1. Start pay transaction

```kotlin
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
}
}.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 ->
service.order(param)
.doOnError { throwable ->
transactionManager.rollback(transactionId, throwable.message)
}
}.doOnSuccess { transactionId ->
transactionManager.commit(transactionId)
}
}
```

#### Scenario3. Check exists transaction

```kotlin
fun exists(param: Any): Mono<Any> {
return transactionManager.exists(param.transactionId) // Find any transaction has ever been started
}
```

#### Scenario4. Handle transaction event

다른 분산서버가 (혹은 자기자신이) transactionManager를 통해서 트랜잭션을 시작하거나 트랜잭션 상태를 변경했을때, 호출한 메소드에 맞는 트랜잭션 이벤트를 발행합니다.

```kotlin

@EventListener(TransactionStartEvent::class)
fun handleTransactionStartEvent(event: TransactionStartEvent) {
// ...
}

@EventListener(TransactionJoinEvent::class)
fun handleTransactionJoinEvent(event: TransactionJoinEvent) {
// ...
}

@EventListener(TransactionCommitEvent::class)
fun handleTransactionCommitEvent(event: TransactionCommitEvent) {
// ...
}

@EventListener(TransactionRollbackEvent::class)
fun handleTransactionRollbackEvent(event: TransactionRollbackEvent) {
// ...
}
```

## Download

```groovy

repositories {
maven { url "https://jitpack.io" }
}

dependencies {
implementation "com.github.rooftop-msa:netx:${version}"
}
```
19 changes: 18 additions & 1 deletion build.gradle
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
import org.springframework.boot.gradle.plugin.SpringBootPlugin

plugins {
id "application"
id "org.jetbrains.kotlin.jvm" version "${jetbrainKotlinVersion}"
id "org.jetbrains.kotlin.plugin.spring" version "${jetbrainKotlinVersion}"
id "org.springframework.boot" version "${springbootVersion}"
id "org.springframework.boot" version "${springbootVersion}" apply false
id "io.spring.dependency-management" version "${springDependencyManagementVersion}"
id "org.sonarqube" version "${sonarcloudVersion}"
id "com.google.protobuf" version "${protobufPluginVersion}"
id "maven-publish"
}

group = "${group}"
Expand All @@ -15,6 +18,20 @@ repositories {
mavenCentral()
}

dependencyManagement {
imports {
mavenBom SpringBootPlugin.BOM_COORDINATES
}
}

publishing {
publications {
maven(MavenPublication) {
from components.java
}
}
}

apply from: "gradle/mq.gradle"
apply from: "gradle/test.gradle"
apply from: "gradle/core.gradle"
Expand Down
7 changes: 2 additions & 5 deletions gradle/spring.gradle
Original file line number Diff line number Diff line change
@@ -1,9 +1,6 @@
jar {
enabled = false
}

dependencies {
implementation "org.springframework.boot:spring-boot-starter"
implementation 'org.springframework:spring-context-support'
implementation 'org.springframework.boot:spring-boot-autoconfigure'
implementation "org.springframework.boot:spring-boot-starter-data-redis-reactive"

testImplementation "org.springframework.boot:spring-boot-starter-test"
Expand Down
2 changes: 1 addition & 1 deletion idl
Empty file.
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package org.rooftop.netx.api

data class TransactionCommitEvent(
val transactionId: String,
val nodeName: String,
)
6 changes: 6 additions & 0 deletions src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package org.rooftop.netx.api

data class TransactionJoinEvent(
val transactionId: String,
val nodeName: String,
)
17 changes: 17 additions & 0 deletions src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
package org.rooftop.netx.api

import reactor.core.publisher.Mono

interface TransactionManager {

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

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

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

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

fun rollback(transactionId: String, cause: String): Mono<String>

}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package org.rooftop.netx.api

data class TransactionRollbackEvent(
val transactionId: String,
val replay: String,
val nodeName: String,
val cause: String?,
)
6 changes: 6 additions & 0 deletions src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package org.rooftop.netx.api

data class TransactionStartEvent(
val transactionId: String,
val nodeName: String,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
package org.rooftop.netx.engine

import org.rooftop.netx.api.TransactionCommitEvent
import org.rooftop.netx.api.TransactionJoinEvent
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 reactor.core.publisher.Flux
import reactor.core.publisher.Mono

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

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

protected abstract fun receive(event: SubscribeTransactionEvent): Flux<Transaction>

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}\"")
}
}
}

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

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

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

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