반응형
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 구현체와 통합 가능
반응형
'개발 > Reactive Streams' 카테고리의 다른 글
| Java 에서 비동기 처리를 하는 방법 (Future, Completable Future) (1) | 2025.04.19 |
|---|---|
| [Reactive Streams] Netflix 에서 RxJava 를 도입한 이유 (1) | 2025.04.12 |
