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

Reactive Programming

목차

0. Reactive System(리액티브 시스템)

1. RxJava 개념

2. Reactive Streams

3. Marble Diagram(마블 다이어그램)

4. 비동기 처리

5. 예외 처리

7. 리소스 관리

8. Flowable과 Observable을 생성하는 연산자

9. 통지 데이터를 변환하는 연산자들

학습 흐름

위 링크는 리액티브 시스템의 배경부터 RxJava 연산자까지 이어지는 목차다. 먼저 리액티브 시스템이 요구하는 응답성·탄력성의 배경을 보고, Reactive Streams의 Publisher와 Subscriber 계약을 이해하면 뒤의 배압 설명이 읽기 쉬워진다.

flowchart LR
  시스템[리액티브 시스템] --> 스트림[Reactive Streams 계약]
  스트림 --> rx[RxJava의 Flowable / Observable]
  rx --> 운영[비동기 처리·예외·리소스]
  운영 --> 연산자[생성·변환 연산자]

여기서 리액티브 시스템은 시스템 전체의 설계 성질을, Reactive Streams는 비동기 스트림의 인터페이스와 배압 계약을, RxJava는 스트림을 조합하는 라이브러리를 가리킨다. 이름이 비슷하지만 범위가 다르다. 예제를 볼 때는 누가 값을 만들고, 누가 소비하며, 언제 구독을 해제하는지 함께 따라가면 된다.

주문 이벤트를 처리한다면

새 주문 이벤트가 들어올 때마다 재고를 조회하고 결제를 요청하는 작업을 생각해 보자. 동기 코드에서는 한 요청을 끝까지 기다린 다음 다른 요청을 처리하기 쉽다. 이벤트 스트림으로 표현하면 주문 ID의 흐름에 조회와 변환 연산자를 연결할 수 있지만, Flowable로 감쌌다는 사실만으로 응답 시간이나 장애 복구가 보장되지는 않는다. 재고 API가 느려졌을 때 대기 중인 주문이 몇 개까지 늘어날지, 결제 실패 후 무엇을 다시 시도할지가 별도 설계 사항이다.

import io.reactivex.Flowable;

public class OrderFlow {
    public static void main(String[] args) {
        Flowable.just("order-1", "order-2")
                .map(id -> "처리 대상: " + id)
                .subscribe(System.out::println,
                        error -> System.err.println("실패: " + error.getMessage()),
                        () -> System.out.println("입력 완료"));
    }
}

이 작은 예제는 동기적으로 두 문자열을 바꿔 출력한다. 별도 스레드도, 외부 API도 사용하지 않는다. Flowable을 만들고 subscribe로 소비자를 연결해야 실행 결과가 내려온다는 점과, 값·오류·완료가 서로 다른 신호라는 점만 보여 준다. 운영 코드에서 외부 API를 연결할 때는 호출이 블로킹인지 먼저 확인하고, 필요하면 subscribeOn으로 해당 작업을 적합한 스케줄러에 배치한다. map만 붙인다고 블로킹 호출이 비동기로 바뀌지는 않는다.

같은 단어를 서로 다른 범위에서 쓰는 이유

범위주로 묻는 질문이 묶음의 글
리액티브 시스템부하와 장애에도 응답을 유지하도록 경계와 메시지를 어떻게 설계할까?Reactive System
Reactive Streams소비자가 감당할 만큼만 onNext를 받도록 어떤 계약을 지킬까?Reactive Streams, 배압
RxJava값을 언제 만들고 어떤 순서·스레드·오류 경로로 전달할까?RxJava 개념과 연산자

시스템 설계의 응답성, 회복성, 탄력성, 메시지 기반이라는 성질은 라이브러리 한 개로 완성되지 않는다. Reactive Streams는 수요 요청과 취소를 표준화하지만 외부 API의 장애나 데이터베이스 거래까지 해결하지 않는다. RxJava는 그 위에서 비동기 작업을 조합할 도구를 제공한다. 따라서 이 목차는 먼저 왜 이벤트 흐름을 쓰는지, 다음으로 어떤 계약으로 흐름을 제한하는지, 마지막으로 어떤 연산자를 쓸지 순서로 배열했다.

운영 관점에서 읽기

연산자 예제의 출력만 보지 말고 구독 시점을 먼저 표시해 보자. just(expensiveCall())은 구독 전에 호출이 끝나지만 fromCallable(() -> expensiveCall())은 구독 때 호출된다. flatMap으로 여러 내부 작업을 병합하면 완료 순서에 따라 결과가 섞일 수 있고 concatMap은 순서를 지키는 대신 뒤 작업을 기다리게 한다. 시간 소스인 interval은 소비 속도에 맞춰 무한히 늦출 수 있는 소스가 아니므로 배압과 종료 조건을 함께 읽어야 한다.

실제 처리 흐름을 설계할 때는 입력률, 구독별 처리 지연, 큐 길이, 오류율을 본다. 예를 들어 초당 100건이 들어오는데 소비가 80건이면 1분에 1,200건이 밀린다. buffer로 목록을 묶어도 처리량이 늘지 않으면 이 차이는 사라지지 않는다. 작업 동시성의 상한을 정하거나 생산을 늦추고, 느려진 외부 서비스에는 타임아웃과 실패 경로를 정해야 한다. cancel은 더 이상 필요하지 않은 구독을 정리하는 수단이지만 이미 외부에 반영된 부작용까지 취소하는 거래 명령은 아니다.

참고: Reactive Manifesto, Reactive Streams JVM 명세, RxJava 2 Flowable API.

같은 카테고리의 글