| Spring WebFlux | Reactive Extension Java

아야하면우유·2026년 6월 19일

WebFlux

목록 보기
2/3
post-thumbnail

📄 Reactive Programming
데이터 스트림(Data Stream)과 변경 전파에 관련된 선언적 프로그래밍
패러다임.

반응형 프로그래밍의 키워드는 다음과 같습니다.

  • 비동기성(Asynchrony)
    반응형 시스템은 이벤트 또는 데이터 스트림(Data Stream)을 예상할 수 없는 시간에 처리합니다.
  • 반응성(Responsiveness)
    반응형 시스템은 실시간으로 데이터의 변화에 반응합니다.
  • 탄력성(Elasticity)
    반응형 시스템은 부하나 실패에 유연하게 대응할 수 있습니다.
  • 메시지 기반(based Messaging)
    반응형 시스템은 메시지 기반 아키텍처를 기반으로 동작합니다.

Observables

데이터 흐름에 맞게 알림을 보내 Observer가 데이터를 사용할 수 있도록 한다. Observable은 데이터 스트림을 정의하고, Observer는 이를 구독하여 변환된 데이터를 비동기적으로 처리한다.

💡 Collections의 Iterable과 유사하다

Iterable - 소비자(Consumer)가 값을 요청하고, 준비될 때까지 Thread를 Blocking

Observable - Non-Blocking, 데이터가 준비되면 Consumer에게 Push

메시지 큐잉에서 주로 사용하는 Publisher인듯 하다...!


Consume 방식

Observable이 데이터를 발행한 후, 알림(Event) 을 전파하면, Observer는 그것을 구독(Subscribe) 하고, 데이터를 소비(Consume) 한다.

Observable의 데이터 발행

Observable이 데이터를 발행한 후, 보내는 알림에는 세 종류가 있다.

public interface Emitter(@NonNull T> {
	void onNext(@NonNull T value);
    void onError(@NonNull Throwable error);
    void onComplete();
}
  • onNext : 데이터의 발행을 알린다.
  • onComplete : 모든 데이터의 발행이 완료됨을 알린다. 이후 onNext 를 호출하면 안 된다.
  • onError : 오류가 발생했음을 알린다. 이후에 onNext, onComplete 는 호출되지 않는다.

Observable의 Subscribe

구독(Subscribe) 이란, 데이터 발행을 수신받고, 해당 데이터를 활용하여 다른 작업을 진행하기 위한 행위를 의미한다. Observer는 subscribe() 메소드에서 수신한 각 알림에 대해 실행할 내용을 선언한다.

  • Disposable
    사전적 의미 : 사용 후 버리는, 일회용의 / 이용 가능한
public interface Disposable {
	void dispose()	// 리소스를 해제, 해제 시 작업은 멱등성을 가짐.
    boolean isDisposed()	// 리소스 해제 여부를 boolean으로 반환
}

Observersubscribe() 메소드의 반환 타입인 Disposable 인터페이스에 대해 짚은 후에 subscribe에 대해 더 알아보는 것이 좋을 것 같아, Disposable 인터페이스에 대한 설명을 작성했다.

Disposable은 Observer는 Observable을 구독하고, 구독함으로써 스트림을 생성한다. 스트림이 오래 실행될 경우, 메모리 누수가 발생하므로 이 스트림을 정리해야 할 필요성이 있다.

import io.reactivex.rxjava3.core.Observable;
import io.reactivex.rxjava3.disposables.Disposable;

import java.util.concurrent.TimeUnit;

public class Main {

	public static void main(String[] args) {
    	
        Observable<Long> observable = Observable.interval(1, TimeUnit.SECONDS);
        Disposable disposable = observable.subscribe(System.out::println);
        
        new Thread(() -> {
        	try {
            	Thread.sleep(2500);
            } catch (Exception e) {
            	e.printStackTrace();
            }
            
            System.out.println("isDisposed: " + disposable.isDisposed());
            System.out.println("called Disposable.dispose()...");
            disposable.dispose();
            System.out.println("isDisposed: " + disposable.isDisposed());
        }).start();
    }
}

Observable은 1초마다 1을 발행한다.
Observer는 이를 구독하고, 람다식을 통해 콘솔에 발행된 아이템을 출력한다.

이 행위는 Thread.sleep(2500)을 통해 2.5초간 실행된다.

  • 0 ~ 1 = 0
  • 1 ~ 2 = 1
  • 2 ~ 2.5 = null

이후 Thread가 종료되고, 이전에 subscribe()를 사용해 구독했던 observer의 구독 리소스 해제를 진행한다.

이제 Disposable의 사용법을 대충 알아보았으니, subscribe()의 반환 형태를 알아보자.

public final Disposable subscribe()
public final Disposable subscribe(@NonNull Consumer<? super T> on Next)
public final Disposable subscribe(@NonNull Consumer<? super T> onNext, @NonNull(? super Throwable) onError)
public final Disposable subscribe(@NonNull Consumer<? super T> onNext, @NonNull(? super Throwable) onError, @NonNull Action onComplete)
public final void subscribe(@NonNull Observer<? super T> observer)

인자가 없는 subscribe()의 경우, 테스트나 디버깅할 때 주로 사용하며, onError 이벤트가 발생하면 onErrorNotImplementedException을 던진다.

💡 onComplete 이벤트가 발생하면, dispose()를 호출하고
Observable이 더 이상 데이터를 발행하지 않게 구독을 해지한다.


Operators

RxJava는 다양한 연산자를 제공하여 Observable의 데이터를 변환, 조작, 필터링, 결합 등의 작업을 수행한다.

  • just() : 단일 항목을 emit
  • listOf() : Iterable 또는 Future에서 Obsavable을 생성
  • interval() : 주기적으로 emit
  • timer() : 지정된 시간 이후에 emit

등등 여러 타입의 Observable을 생성한다.


Scheduler

RxJava는 스케줄러를 통해 비동기 작업을 관리하고, 스레드 풀을 활용하여 병렬 처리를 지원합니다. 이를 통해 I/O 작업, 네트워크 호출, 데이터베이스 쿼리 등 비동기 작업을 효과적으로 처리할 수 있습니다.

profile
우유가 넘어지면 아야

0개의 댓글