📄 Reactive Programming
데이터 스트림(Data Stream)과 변경 전파에 관련된 선언적 프로그래밍
패러다임.
반응형 프로그래밍의 키워드는 다음과 같습니다.
데이터 흐름에 맞게 알림을 보내 Observer가 데이터를 사용할 수 있도록 한다. Observable은 데이터 스트림을 정의하고, Observer는 이를 구독하여 변환된 데이터를 비동기적으로 처리한다.
💡 Collections의 Iterable과 유사하다
Iterable - 소비자(Consumer)가 값을 요청하고, 준비될 때까지 Thread를 Blocking
Observable - Non-Blocking, 데이터가 준비되면 Consumer에게 Push
메시지 큐잉에서 주로 사용하는 Publisher인듯 하다...!
Observable이 데이터를 발행한 후, 알림(Event) 을 전파하면, Observer는 그것을 구독(Subscribe) 하고, 데이터를 소비(Consume) 한다.
Observable이 데이터를 발행한 후, 보내는 알림에는 세 종류가 있다.
public interface Emitter(@NonNull T> {
void onNext(@NonNull T value);
void onError(@NonNull Throwable error);
void onComplete();
}
onNext : 데이터의 발행을 알린다.onComplete : 모든 데이터의 발행이 완료됨을 알린다. 이후 onNext 를 호출하면 안 된다.onError : 오류가 발생했음을 알린다. 이후에 onNext, onComplete 는 호출되지 않는다.구독(Subscribe) 이란, 데이터 발행을 수신받고, 해당 데이터를 활용하여 다른 작업을 진행하기 위한 행위를 의미한다. Observer는 subscribe() 메소드에서 수신한 각 알림에 대해 실행할 내용을 선언한다.
Disposablepublic interface Disposable {
void dispose() // 리소스를 해제, 해제 시 작업은 멱등성을 가짐.
boolean isDisposed() // 리소스 해제 여부를 boolean으로 반환
}
Observer의 subscribe() 메소드의 반환 타입인 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초간 실행된다.
이후 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이 더 이상 데이터를 발행하지 않게 구독을 해지한다.
RxJava는 다양한 연산자를 제공하여 Observable의 데이터를 변환, 조작, 필터링, 결합 등의 작업을 수행한다.
just() : 단일 항목을 emitlistOf() : Iterable 또는 Future에서 Obsavable을 생성interval() : 주기적으로 emittimer() : 지정된 시간 이후에 emit등등 여러 타입의 Observable을 생성한다.
RxJava는 스케줄러를 통해 비동기 작업을 관리하고, 스레드 풀을 활용하여 병렬 처리를 지원합니다. 이를 통해 I/O 작업, 네트워크 호출, 데이터베이스 쿼리 등 비동기 작업을 효과적으로 처리할 수 있습니다.