Reactor와 마블 다이어그램

...·2024년 9월 27일

Reactive Programming

목록 보기
5/5
post-thumbnail

Reactor란?

Reactor는 Spring Framework 팀의 주도하에 개발된 리액티브 스트림즈의 구현체로서 Spring Framework 5 버전부터 리액티브 스택에 포함되어 Spring WebFlux 기반의 리액티브 애플리케이션을 제작하기 위한 핵심 역할을 담당한다.

리액티브 스트림즈의 구현체인 Reactor는 리액티브 프로그래밍을 위한 라이브러리라고 정의할 수 있다. 따라서 Reactor Core 라이브러리는 Spring WebFlux 프레임워크에 라이브러리로 포함되어 있다.

  1. Reactive Streams: Reactor는 리액티브 스트림즈 사양을 구현한 리액티브 라이브러리이다.
  2. Non-Blocking: Reactor는 JVM 위에서 실행되는 Non-Blocking 애플리케이션을 제작하기 위해 필요한 핵심 기술이다.
  3. Java's functional API: Reactor에서 Publisher와 Subscriber 간의 상호작용은 Java의 함수형 프로그래밍 API를 통해서 이루어진다.
  4. Flux[N]: Reactor의 Publisher 타입은 크게 두 가지인데, 그중 한 가지가 바로 Flux이다. Flux[N]이라는 의미는 N개의 데이터를 emit한다는 것인데, 다시 말해서 Flux는 0개부터 N개, 즉 무한대의 데이터를 emit할 수 있는 Reactor의 Publisher이다.
  5. Mono[0|1]: Mono 역시 Reactor에서 지원하는 Publisher 타입인데, Mono[0|1]과 같이 표현된 이유는 Mono가 데이터를 한 건도 emit하지 않거나 단 한 건만 emit하는 단발성 데이터 emit에 특화된 Publisher이기 때문이다.
  6. Well-suited for microservices: Non-Blocking I/O 특징을 가지는 Reactor는 마이크로 서비스 기반 시스템에서 수많은 서비스들 간에 지속적으로 발생하는 I/O를 처리하기에 적합한 기술이다.
  7. Backpressure-ready network: Reactor는 Publisher로부터 전달받은 데이터를 처리하는 데 있어 과부하가 걸리지 않도록 제어하는 Backpressure를 지원한다.

Hello Reactor 코드로 보는 Reactor의 구성요소

public class Example {
	public static void main(String[] args) {
    	Flux<String> sequence = Flux.just("Hello", "Reactor");
        sequence.map(data -> data.toLowerCase())
        	.subscribe(data -> System.out.println(data));
    }
}

Flux는 Reactor에서 Publisher의 역할을 한다. 입력으로 들어오는 데이터는 just 메서드의 파라미터로 전달한 "Hello", "Reactor"이다. 이 데이터는 Publisher가 최초로 제공하는 가공되지 않은 데이터로서 데이터 소스라고 불린다.

이 데이터 소스의 데이터 개수가 둘이기 때문에 N건의 데이터를 처리할 수 있는 Reactor의 Publisher 타입인 Flux를 사용했다.

subscribe 메서드의 파라미터로 전달된 람다 표현식인 data -> System.out.println(data)가 바로 Subscriber 역할을 한다. 이 람다 표현식은 Consumer 함수형 인터페이스이며, 내부적으로는 LambdaSubscriber라는 클래스에 전달되어 데이터를 처리하는 역할을 하게 된다.

just(), map()은 Reactor에서 지원하는 Operator 메서드인데, just() Operator는 데이터를 생성해서 제공하는 역할을 하고, map() Operator는 전달받은 데이터를 가공하는 역할을 한다. map() Operator의 파라미터로 정의된 람다 표현식에서는 just() Operator로부터 전달받은 문자열 데이터를 소문자로 변경한다.

그리고 just() Operator의 리턴 값이 Flux임을 확인할 수 있다. 이를 통해 Reactor의 Operator는 리턴 값으로 Flux(또는 Mono)를 반환하기 때문에 Operator 체인을 형성하여 또 다른 Operator를 연속적으로 호출할 수 있다는 사실을 알 수 있다.

Reactor의 Hello, World 코드는 단순하지만 Reactor의 핵심 구성요소를 대부분 포함한다. '데이터를 생성해서 제공하고(1단계), 데이터를 가공한 후에(2단계) 전달받은 데이터를 처리한다(3단계)'는 이 세 가지 단계는 데이터를 가공 처리하는 단계가 얼마나 복잡해지느냐, 스레드를 어떤 식으로 제어하느냐 등의 추가 작업과 상관없이 수행되는 필수 단계라는 사실을 기억해야 한다.

마블 다이어그램(Marble Diagram)이란?

  1. 다이어그램에는 두 개의 타임라인이 존재하는데, 첫째가 바로 Publisher가 데이터를 emit하는 타임라인이다. 이 Publisher는 데이터 소스를 최초로 emit하는 Publisher일 수도 있고 그렇지 않을 수도 있다. Operator 함수를 기준으로 상위에 있는, 즉 Upstream의 Publisher라고 보는 것이 적절하다. (Reactor에서는 Flux의 경우 Source Flux라고도 부른다.)
  2. (2)는 Publisher가 emit하는 데이터를 의미한다. 타임라인은 왼쪽에서 오른쪽으로 시간이 흐르는 것을 의미하기 때문에 가장 왼쪽에 있는 1번 구슬이 시간상으로 가장 먼저 emit된 데이터이다.
  3. (3)의 수직으로 된 바는 데이터의 emit이 정상적으로 끝났음을 의미한다.
  4. Operator 함수 쪽으로 들어가는 점선 화살표는 Publisher로부터 emit된 데이터가 Operator 함수 쪽으로 들어가는 점선 화살표는 Publisher로부터 emit된 데이터가 Operator 함수의 입력으로 전달되는 것을 의미한다.
  5. Publisher로부터 전달받은 데이터를 처리하는 Operator 함수이다.Reactor는 굉장히 많은 수의 Operator를 지원한다.
  6. Operator 함수에서 나가는 점선 화살표는 Publisher로부터 전달받은 데이터를 가공 처리한 후에 출력으로 내보내는 것을 의미한다. 출력으로 내보낸다는 의미는 정확하게 표현하자면, Operator 함수에서 리턴하는 새로운 Publisher를 이용해 Downstream에 가공된 데이터를 전달하는 것을 의미힌다.
  7. Operator 함수에서 가공 처리되어 출력으로 내보내진 데이터의 타임라인이다.(Reactor에서는 Operator의 출력으로 리턴된 Flux의 경우, Output Flux라고도 부른다.)
  8. 'X' 표시는 에러가 발생해 데이터 처리가 종료되었음을 의미하며, onError Signal에 해당된다.

처음 사용해 보는 Operator를 올바르게 이해하고 사용하기 위해서 해당 Operator의 API 설명을 보기 전에 마블 다이어그램부터 먼저 확인하는 습관을 들이는 것이 좋다.

마블 다이어그램으로 Reactor의 Publisher 이해하기

Mono는 단 하나의 데이터를 emit하는 Publisher이기 때문에 그림에서도 단 하나의 데이터만 표현한다.

public class Example {
	public static void main(String[] args) {
    	Mono.just("Hello Reactor")
        	.subscribe(System.out::println);
    }
}

just() Operator는 한 개 이상의 데이터를 emit하기 위한 대표적인 Operator로서 2개 이상의 데이터를 파라미터로 전달할 경우, 내부적으로 fromArray() Operator를 이용해서 데이터를 emit한다.

public class Example {
	public static void main(String[] args) {
    	Mono
        	.empty()
            .subscribe(
            	none -> System.out.println("# emitted onNext signal"),
                error -> {},
                () -> System.out.println("# emitted onComplete signal")
            );
    }
}

데이터를 0건, 즉 데이터를 한 건도 emit하지 않는 예제 코드이다. empty() Operator를 사용하면 데이터를 emit하지 않고 onComplete Signal을 전송한다.

이처럼 데이터를 한 건도 emit하지 않는 empty() Operator는 주로 어떤 특정 작업을 통해 데이터를 전달받을 필요는 없지만 작업이 끝났음을 알리고 이에 따른 후처리를 하고 싶을 때 사용할 수 있다.

Flux 활용 예제

public class Example {
	public static void main(String[] args) {
    	Flux<String> flux = 
        	Mono.justOrEmpty("Steve")
            	.concatWith(Mono.justOrEmpty("Jobs"));
        flux.subscribe(System.out::println);
    }
}

두 개의 Mono를 연결해서 Flux로 변환했다.

just() Operator의 경우 파라미터의 값으로 null을 허용하지 않지만 justOrEmpty()는 null을 허용한다. justOrEmpty()의 파라미터로 null이 전달되면 내부적으로 empty() Operator를 호출하도록 구현되어 있다.

concatWith() Operator는 concatWith()를 호출하는 Publisher와 concatWith()의 파라미터로 전달되는 Publisher가 각각 emit하는 데이터들을 하나로 연결해서 새로운 Publisher의 데이터 소스로 만들어 주는 Operator이다.

마블 다이어그램을 보면 concatWith() 위쪽에 있는 Publisher의 데이터 소스와 concatWith 내부에 있는 Publisher의 데이터 소스를 연결하는 것을 볼 수 있다.

이렇게 연결된 데이터 소스는 새로운 Flux의 데이터 소스가 되어 차례대로 emit된다.

public class Example {
	public static void main(String[] args) {
    	Flux.concat(
        			Flux.just("Mercury", "Venus", "Earth"),
                    Flux.just("Mars", "Jupiter", "Saturn"),
                    Flux.just("Uranus", "Neptune", "Pluto")
                .collectList()
                .subscribe(planets -> System.out.println(planets));
    }
}

concatWith()의 경우 두 개의 데이터 소스만 연결할 수 있지만, concat()은 여러 개의 데이터 소스를 원하는 만큼 연결할 수 있다.

collectList() Operator는 Upstream Publisher에서 emit하는 데이터를 모아서 List의 원소로 포함시킨 새로운 데이터 소스로 만들어 주는 Operator이다.

  • concat() Operator에서 리턴하는 Publisher는 Flux이다. Mono는 0개 또는 1개의 데이터만 emit할 수 있는 Publisher이고, Flux는 N개, 즉 여러 건의 데이터를 emit할 수 있는 Publisher이다. concat() Operator를 이용해서 총 세 개의 Flux가 가지고 있는 아홉 개의 데이터를 데이터 소스로 연결하기 때문에 concat() Operator의 리턴 값은 Flux가 될 수밖에 없다.
  • collectList() Operator는 여러 개의 데이터를 하나의 List에 원소로 포함시킨다. 즉, List에 포함된 원소는 여러 개이지만 List 자체는 하나이기 때문에 한 개의 데이터만 emit할 수 있는 Mono를 리턴한다.

참고 서적: 스프링으로 시작하는 리액티브 프로그래밍

profile
주니어 백엔드 개발자

0개의 댓글