본문 바로가기

Reactor 의 Mono, Flux 는 Reactive Streams 의 구현체이다

@whateverU2025. 5. 14. 23:31
반응형

Reactor 의 Mono, Flux 는 Reactive Streams 의 구현체이다.
Reactive Streams 인터페이스와 Reactor 를 직접 구현하면서 알아보자.

Reactive Streams 구현체 예제

의존성

dependencies {
    testImplementation(kotlin("test"))

    // https://mvnrepository.com/artifact/org.reactivestreams/reactive-streams
    implementation("org.reactivestreams:reactive-streams:1.0.4")

    // https://mvnrepository.com/artifact/org.junit.jupiter/junit-jupiter-api
    implementation("org.junit.jupiter:junit-jupiter-api:5.12.1")

    // https://mvnrepository.com/artifact/io.projectreactor/reactor-core
    implementation("io.projectreactor:reactor-core:3.8.0-M2")
}

Publisher

class SimplePublisher(private val data: List<Int>) : Publisher<Int> {
    override fun subscribe(subscriber: Subscriber<in Int>) {
        val subscription = AsyncSubscription(subscriber, data)
        subscriber.onSubscribe(subscription)
    }

    private class AsyncSubscription(
        private val subscriber: Subscriber<in Int>,
        private val data: List<Int>
    ) : Subscription {
        private val requested = AtomicLong(0)
        private val cancelled = AtomicBoolean(false)
        private val emitting = AtomicBoolean(false)
        private var index = 0
        private val executor = Executors.newSingleThreadExecutor()

        override fun request(n: Long) {
            if (n <= 0) {
                subscriber.onError(IllegalArgumentException("Request must be > 0"))
                return
            }

            requested.addAndGet(n)
            tryEmit()
        }

        override fun cancel() {
            cancelled.set(true)
            executor.shutdownNow()
        }

        private fun tryEmit() {
            if (emitting.compareAndSet(false, true)) {
                executor.submit {
                    try {
                        while (requested.get() > 0 && index < data.size && !cancelled.get()) {
                            Thread.sleep(300)
                            subscriber.onNext(data[index++])
                            requested.decrementAndGet()
                        }

                        if (index == data.size && !cancelled.get()) {
                            subscriber.onComplete()
                            executor.shutdown()
                        }
                    } catch (e: Exception) {
                        subscriber.onError(e)
                        executor.shutdown()
                    } finally {
                        emitting.set(false)
                        // 남은 요청이 있으면 다시 emit
                        if (requested.get() > 0 && index < data.size && !cancelled.get()) {
                            tryEmit()
                        }
                    }
                }
            }
        }
    }
}

Subscriber

class SimpleSubscriber : Subscriber<Int> {
    private var subscription: Subscription? = null

    override fun onSubscribe(s: Subscription) {
        println("[${Thread.currentThread().name}] Subscribed")
        subscription = s
        s.request(1)
    }

    override fun onNext(t: Int) {
        println("[${Thread.currentThread().name}] Received: $t")
        Thread.sleep(500)
        subscription?.request(1)
    }

    override fun onError(t: Throwable) {
        println("[${Thread.currentThread().name}] Error: ${t.message}")
    }

    override fun onComplete() {
        println("[${Thread.currentThread().name}] Done!")
    }
}

Try - Reactive Streams

main-code

fun main() {
    println("[${Thread.currentThread().name}] main thread start...")
    val publisher = AsyncPublisher((1..10).toList())
    val subscriber = SimpleSubscriber()
    publisher.subscribe(subscriber = subscriber)
}

output

[main] main thread start...
[main] Subscribed
[pool-1-thread-1] Received: 1
[pool-1-thread-1] Received: 2
[pool-1-thread-1] Received: 3
[pool-1-thread-1] Received: 4
[pool-1-thread-1] Received: 5
[pool-1-thread-1] Received: 6
[pool-1-thread-1] Received: 7
[pool-1-thread-1] Received: 8
[pool-1-thread-1] Received: 9
[pool-1-thread-1] Received: 10
[pool-1-thread-1] Done!

Try - Reactor

main-code

Reactor Flux 를 생성하고 Reactive Streams 구현한 SimpleSubscriber 가 subscribe 하면? 된다? 왜? Flux 또한 Publisher 와 동일하게 Reactive Streams를 구현한 구현체이기 때문에 가능하다.

fun main(){
    val executor = Executors.newSingleThreadExecutor()
    val scheduler = Schedulers.fromExecutor(executor)

    println("[${Thread.currentThread().name}] main thread start...")
    val simpledSubscriber = SimpleSubscriber()
    Flux.just(1, 2, 3, 4, 5)
        .subscribeOn(scheduler)
        .doFinally {
            executor.shutdown()
        }
        .subscribe(
            simpledSubscriber
        )
}

output

[main] main thread start...
[main] Subscribed
[pool-1-thread-1] Received: 1
[pool-1-thread-1] Received: 2
[pool-1-thread-1] Received: 3
[pool-1-thread-1] Received: 4
[pool-1-thread-1] Received: 5
[pool-1-thread-1] Done!

결론

이게 되네?

Reactive Streams 구현체를 만들어 보면 알겠지만 인터페이스는 정말 단순 한데 구현하기 쉽지 않음. https://github.com/reactive-streams/reactive-streams-jvm?tab=readme-ov-file#specification 단순한 인터페이스에 비해 룰이 겁나게 많다

그래서 똑똑한 분들이 자 쉽죠? 하고 만든 여러 라이브러리. 우리는 이해하고 쓰면 된다.
왜 만들었는지? 어떻게 동작하는지?

ReactiveStreams 스팩을 구현한 라이브러리 (By AI)

  • Project Reactor
    • Spring WebFlux의 핵심 구현체
    • Mono, Flux 타입 제공
    • 비동기, 논블로킹 스트림 처리
  • RxJava
    • ReactiveX 프로젝트의 Java 구현
    • Observable, Flowable 등 다양한 타입 제공
    • Reactor보다는 범용적이고 연산자가 더 많음
  • Akka Streams
    • Actor 기반의 Reactive Stream 처리
    • 높은 안정성과 백프레셔 지원
    • 복잡하지만 강력한 도구
  • Mutiny
    • Quarkus에서 주로 사용
    • Uni, Multi 타입 제공
    • 간결하고 직관적인 API
  • Vert.x Reactive Streams
    • Vert.x 기반 환경에 최적화
    • 다른 Reactive Streams 구현체와 통합 가능
반응형
whateverU
@whateverU :: whateverU

sang12.co.kr https://github.com/ChoiSangIl

공감하셨다면 ❤️ 구독도 환영합니다! 🤗

목차