본문으로 건너뛰기
홈
기술
기술 전체
프로그래밍68
컴퓨터 과학63
AI48
웹 개발36
인프라33
데이터31
소프트웨어 공학18
소개
← 목록으로프로그래밍 › 비동기 › 리액티브 프로그래밍

1. RxJava 개념

목차

RxJava 개념

RxJava는 Reactive Programming을 구현하는데 사용하는 라이브러리 이다. 에릭 마이어가 개발한 .NET 프레임워크의 실험적 라이브러리인 Reactive Extensions(Rx) 를 2009년 마이크로소프트에서 공개하고 2013년 넷플릭스에서 자바로 이식한 것이 RxJava의 시작이다.

현재 Reactive Extensions를 다루는 라이브러리는 ReactiveX라는 오픈 소스 프로젝트로 바뀌어 자바와 .NET 뿐만 아니라 자바스크립트, 스위프트, 등 여러 프로그램언어를 지원하는 라이브러리를 제공한다.

Reactive Extensions는 동기식 또는 비동기식 스트림과 관계 없이 명령형 언어를 이용해 데이터 스트림을 조작할 수 있는 일련의 도구이다.

ReactiveX는 Observer Pattern, Iterator Pattern 및 함수형 프로그램의 조합으로 정의된다.

ReactiveX : http://reactivex.io

ReactiveX : http://reactivex.io

특징

  • RxJava는 Observer Pattern을 활용하였다. Observer 패턴은 감시 대상 객체의 상태가 변화면 이를 관찰하는 객체에 알려주는 구조이다. 그래서 데이터를 생성하는 측과 데이터를 소비하는 측으로 나눌 수 있고 쉽게 데이터 스트림을 처리할 수 있다.

  • 스케줄러를 지정하면 비동기 경계에서 작업할 수 있다. 연산자를 연결하기만 하면 자동으로 별도 스레드에서 실행되는 것은 아니다. Reactive Streams 배압 계약은 RxJava 2의 Flowable에 적용되며 Observable의 Observer는 request(n)을 사용하지 않는다.

기본구조

RxJava는 데이터를 만들고 통지하는 생산자(Publisher)와 통지된 데이터를 받아 처리하는 소비자(Subscriber)로 구성된다. 소비자는 생산자를 구독하여 생산자가 통지한 데이터를 소비자가 받아서 처리한다.

크게 두 가지 로 나누는데, Reactive Streams을 지원하는 Flowable과 Subscriber가 있고, Reactive Streams을 지원하지 않는 Observable과 Observer가 있다.

구분생산자소비자
Reactive Streams 지원FlowableSubscriber
Reactive Streams 미지원ObservableObserver

Flowable로 구독시작(onSubscribe)을 하면 데이터 통지(onNext), 에러통지 (onError), 완료통지(onComplete)를 수행하고 통지받은 시점의 소비자인 Subscriber로 처리하고 데이터 개수 요청 및 구독 해지를 할 수 있다. RxJava 2의 Flowable은 Reactive Streams Publisher를 구현한다.

Observable과 Observer Flowable과 onSubScriber와 같은 기능을 수행하지만 통지하는 데이터 개수를 제어하는 배압 기능이 없기 때문에 데이터 개수를 요청하지 않는다. 그래서 onSubScription을 사용하지 않고 Disposable이라는 구독 해지 메서드가 있는 인터페이스를 사용한다. Disposable은 구독 시작 시점에 onSubscribe 메서드의 인자로 Observer에게 전달된다.

Disposable는 구독 해지를 위한 두 가지 메서드가 있다.

methoddescription
dispose구독을 해지한다.
isDisposed구독을 하지하면 true, 하지않으면 false를 반환한다

그러므로 Observable과 Observer은 데이터 개수 요청을 하지 않고 데이터가 생성되자마자 Observer에게 통지된다.

연산자

RxJava는 생산자가 통지한 데이터가 소비자에게 도착하기 전에 불필요한 데이터를 삭제하거나 소비자가 사용하기 쉽게 변경해야 할 때가 있는데 이 때 Flowable/Observable의 메서드에서 새로운 Flowable/Observable을 반환하며 해당 메서드를 서로 연결해나가며 최동 데이터를 통지하는 Flowable/Observable을 생성한다. 그래서 통지하는 데이터를 생성하거나 필터링 또는 변환하는 메서드이다.

rx=operator-info.jpeg

연산자를 사용하여 데이터를 출력하는 예제

[Subscription.java]

import io.reactivex.Flowable;
import io.reactivex.functions.Consumer;
import io.reactivex.functions.Function;
import io.reactivex.functions.Predicate;

public class Subscription {
    public static void main(String[] args) {
        evenRamda();
        evenNotRamda();
    }

    private static void evenRamda() {
        Flowable<Integer> flowable = Flowable.just(1,2,3,4,5,6,7,8,9,10) //인자의 데이터를 순서대로 통지하는 Flowable 생성
                .filter(data -> data % 2 == 0) //짝수 데이터만 통지
                .map(data -> data * 100); //데이터를 100배로 변환

        //구독하여 받은 데이터를 출력
        flowable.subscribe(data-> System.out.println("Data : " + data));
    }

    private static void evenNotRamda() {
        Flowable<Integer> flowable = Flowable.just(1,2,3,4,5,6,7,8,9,10) //인자의 데이터를 순서대로 통지하는 Flowable 생성
                //짝수 데이터만 통지
                .filter(new Predicate<Integer>() {
                    @Override
                    public boolean test(Integer src) throws Exception {
                        return src % 2 == 0;
                    }
                })
                //데이터를 100배로 변환
                .map(new Function<Integer, Integer>() {
                    @Override
                    public Integer apply(Integer data) throws Exception {
                        return data * 100;
                    }
                });

        //구독하여 받은 데이터를 출력
        flowable.subscribe(new Consumer<Integer>() {
            @Override
            public void accept(Integer data) throws Exception {
                System.out.println("Data : " + data);
            }
        });
    }
}

비동기 처리

RxJava는 개발자가 직접 스레드를 관리하지 않게 각 처리 목적에 맞춰 스레드를 관리하는 스케줄러(Scheduler)를 제공하며 이 스케줄러를 이용하면 어떤 스레드에서 무엇을 처리할지 제어할 수 있다.

스케줄러는 데이터를 생성해 통지(flowable/Observable)하는 부분과 데이터를 받아 처리하는 부분을 설정 할 수 있다. 이후 데이터의 필터나 변환을 하는 SubScriber/Observer가 데이터 수신 처리를 어느 스케줄러에서 처리 할 지를 제어한다.

비동기 생성

//1초마다 0부터 시작하는 값을 비동기로 통지하는 interval 생성
Observable.interval(1, TimeUnit.SECONDS)
.subscribe(e -> { //결과 통지 출력
          System.out.println("Received : " + e);
      });
      Thread.sleep(5000);

결과

Received : 0
Received : 1
Received : 2
Received : 3
Received : 4

Thread.sleep(5000);을 지우면 아무것도 출력하지 않고 종료하게 되는데 이벤트가 생성되는 것과 별개의 스레드에서 사용되기 때문이다.

공유 상태를 읽는 비동기 연산의 위험

옛 예제는 static State calcMethod를 scan의 함수에서 읽으면서 다른 스레드의 main이 그 값을 바꿨다. 출력이 중간에 덧셈에서 곱셈으로 바뀐다는 설명은 우연한 실행 순서에 의존한다. static 필드는 구독마다 따로 생기지도 않고, 일반 필드의 스레드 간 변경은 가시성도 보장되지 않는다. 다음 그림은 외부 상태가 데이터 흐름 안으로 끼어드는 지점을 읽는 용도로 남긴다.

비동기 스트림에서 외부 공유 상태를 참조하는 흐름

그림의 외부 화살표처럼 계산 도중 설정이 바뀌면 같은 입력도 구독 시점이나 스레드 스케줄에 따라 다른 결과가 될 수 있다. 한 작업에 하나의 계산 규칙이 필요하다면 구독을 만들 때 지역 값으로 고정한다.

import io.reactivex.Flowable;
import io.reactivex.schedulers.Schedulers;

public class StableCalculation {
    enum Operation { ADD, MULTIPLY }

    public static void main(String[] args) {
        final Operation selected = Operation.ADD;
        Flowable.range(1, 4)
                .subscribeOn(Schedulers.computation())
                .scan(0, (total, value) -> selected == Operation.ADD
                        ? total + value : total * value)
                .blockingForEach(System.out::println);
        // 초기값 0과 각 단계의 누적값: 0, 1, 3, 6, 10
    }
}

subscribeOn은 상위의 range와 누적 계산이 시작되는 스케줄러를 지정한다. blockingForEach는 예제의 main이 완료를 기다리도록 하는 소비 방법이며 비동기 실행 자체를 없애지 않는다. 실제 비동기 처리에서는 구독을 보관하고 종료 시 취소하는 방식을 설계한다. 실행 중 설정 변경이 의도된 요구라면 변경도 스트림의 입력 신호로 모델링하거나, 스레드 안전한 상태 전달 규칙을 정해야 한다.

Cold Constructor와 Hot Constructor

[Cold Constructor]

Cold 소스는 구독할 때마다 해당 구독의 데이터 타임라인을 새로 만든다. 소비자가 한 명으로 제한된다는 뜻은 아니다. Flowable.range(1, 3)을 두 번 구독하면 각 소비자가 1, 2, 3을 받는다.

cold-constructor.jpeg — cold-constructor

[Hot Constructor]

Hot 소스는 생산 타임라인을 여러 구독자가 공유할 수 있다. 이미 발행한 값은 뒤늦게 구독한 소비자가 놓칠 수 있다. replay처럼 과거 값을 저장하는 방식을 붙이면 새 소비자가 저장된 값도 받을 수 있다. 시작 시점은 Hot 여부만으로 정해지지 않는다. 예를 들어 publish()가 만든 ConnectableFlowable은 connect()를 호출해야 원본 구독을 시작한다.

hot-constructor.jpeg

ConnectableFlowable/ConnectableObservable 클래스

ConnectableFlowable/ConnectableObservable는 Hot Flowable/Observable이고 여러 Subscriber/Observer에서 동시에 구독할 수 있다. 또한 subscribe메서드를 호출해도 처리를 시작하지 않고 connect메서드를 호출해야 처리를 한다. 1. subscriber/Observer에서 구독(connect 메서드 호출 전까지 처리되지 않음) 2. connect메서드를 호출 할 때 동시에 여러 구독자에게 데이터를 통지

[Flowable/Observable로 변환하는 메서드]

  • refcount() : 새로운 Flowable/Observable을 생성

  • autoConnect() / autoConnect(int numberOfSubscribers) : 지정한 개수의 구독이 시작된 시점에 처리를 시작하는 Flowable/Observable을 생성

[Flowable/Observable을 Cold에서 Hot으로 변환하는 연산자]

  • publish() : Flowable/Observable에서 ConnectableFlowable/ConnectableObservable을 생성하는 연산자이다. 해당 클래스를 이용하면 처리를 시작한 뒤 구독하면 구독한 이후 생성된 데이터부터 새로운 소비자에게 통지한다.

  • replay() / replay(int bufferSize) / replay(long time, TimeUnit unit) : ConnectableFlowable/ConnectableObservable를 생성하는 연산자로 통지한 데이터를 캐시하고, 처리르 시작한 후 구독하면 캐시된 데이터를 먼저 새로 구독한 소비자에게 통지하며 그 후에 모든 소비자에게 같은 데이터를 통지한다. 그리고 메서드가 없으면 모든 데이터를 캐시하고 인자가 있으면 지정한 시간동안 정한 개수만큼 데이터를 캐시한다.

  • share() : 여러소비자가 구독할 수 있는 Flowable/Observable을 생성한다. 다른메서드와 달리 ConnectableFlowable/ConnectableObservable를 생성하지 않고, Flowable/Observable을 구독하는 소비자가 있는 동안 중간에 새로 구독해도 같은 타임라인에서 생성되는 데이터를 통지한다.

Flowable vs Observable

RxJava 2에서 Flowable은 Reactive Streams Publisher이므로 하위의 요청량을 반영할 수 있다. Observable은 request(n) 배압 계약을 제공하지 않으며 Disposable로 구독을 종료한다. 원본이 느려질 수 있는 파일 읽기나 요청량을 조절할 수 있는 소스라면 Flowable이 맞는다. UI 클릭처럼 이미 발생한 이벤트를 되돌릴 수 없다면 Observable을 사용하고 샘플링·드롭·버퍼 한계 같은 과부하 정책을 따로 정한다.

아래처럼 무한 시간 소스를 만들었다면 타입을 Flowable로 골랐다는 사실만으로 느린 소비자가 안전해지지 않는다. 생산자가 요청량에 맞춰 대기할 수 있는지, 중간의 observeOn 큐가 넘치지 않는지 확인해야 한다. 큰 입력이나 비동기 경계에서 MissingBackpressureException이 발생하는 이유는 Flowable의 배압 계약을 지키려는 과정에서 생산 속도와 소비 속도의 차이를 더 이상 수용할 수 없기 때문이다.

Flowable.interval(1, TimeUnit.MILLISECONDS)
        .take(10)
        .onBackpressureLatest()
        .blockingForEach(System.out::println);

이 코드 조각에는 io.reactivex.Flowable과 java.util.concurrent.TimeUnit import가 필요하다. onBackpressureLatest()는 소비자가 뒤처졌을 때 중간 값을 건너뛰고 최신 값을 유지하는 정책이다. 모든 이벤트를 처리해야 하는 결제·감사 기록에 사용하면 데이터가 사라지므로 적합하지 않다. 초당 몇 건이라는 단일 숫자보다 손실 허용 여부, 입력이 늦춰질 수 있는지, 버퍼 상한과 처리 지연을 기준으로 선택한다.

참고: RxJava 2의 Flowable과 Observable 구분, RxJava 2 배압 설명, RxJava 2 Flowable API.

같은 카테고리의 글