Skip to content

Reactive Programming ‐ Everything About Reactive Programming & Core Concepts of Reactive Programming

woojin edited this page Jul 21, 2026 · 3 revisions

Mono(단일 값 Publisher)

  • 단 하나의 데이터만 발행하거나, 아예 비어있거나, 에러를 내고 종료되는 Publisher이다. 값을 하나 방출하면 즉시 완료된다.

1. 생성(Creation)

  • Mono.just(data) — 이미 계산된 데이터를 감싼다.
  • Mono.empty() — 값이 없음을 표현
  • Mono.defer(Supplier) — 구독 시점까지 실행을 지연
  • Mono.fromCallable(Callable) — 블로킹 작업을 어댑팅

2. 변환(Transformation)

  • map() — 동기 값 변환
  • flatMap() — 비동기 변환 + 평탄화
  • filter() — 조건부 통과

3. 에러 처리(Error Handling)

  • defaultIfEmpty() — 빈 스트림에 기본값 제공
  • switchIfEmpty() — 빈 경우 대체 비동기 흐름으로 전환
  • onErrorReturn() — 에러 시 폴백 값 반환
  • onErrorResume() — 에러를 대체 Mono로 처리

Flux(0~N개 값 Publisher)

  • 0개부터 무한 개까지 값을 방출하고 완료 또는 에러로 종료된다.

1. 생성(Creation)

  • Flux.just() — 여러 개의 사전 정의 값
  • Flux.fromIterable() — 컬렉션을 스트림으로 변환
  • Flux.range() — 숫자 시퀀스 생성
  • Flux.interval() — 시간 기반 방출

2. 변환(Transformation)

  • map() — 원소별 동기 변환
  • flatMap() — 비동기 처리, 순서 보장 안 됨(여러 내부 Publisher 동시 구독)
  • concatMap() — 비동기 순차 처리, 순서 유지
  • collectList()Mono<List<T>>로 집계

주의 : flatMap은 병렬이 아니라 여러 내부 Publisher를 동시에 구독(interleaving)하는 것이다. 기본 동시성 제한(concurrency)이 있으며, 진짜 CPU 병렬은 parallel().runOn(...)을 쓴다.

3. 스트림 에러 처리

  • onErrorResume() — 개별 원소 복구
  • onErrorContinue() — Reactive Streams 계약을 깨는 특수 연산자라 지양. 원소별 복구는 flatMap 안에서 onErrorResume으로 처리하는 게 정석이다.

구독 & Lazy Evaluation

  • 하나 명심할 부분은 아무리 복잡한 Mono/Flux 체인을 만들어도, subscribe()가 호출되기 전엔 아무것도 실행되지 않는다.
  • WebFlux 환경에서는 프레임워크가 구독을 자동 관리한다. 개발자가 직접 subscribe()를 호출할 일은 거의 없다. 컨트롤러가 Mono/Flux를 반환하면 Netty가 내부적으로 구독을 처리한다.
// (함정) just는 인자를 assembly time에 즉시 평가한다 → 구독 전에 실행됨
Mono.just(blockingCall());                // ❌ eager
Mono.fromCallable(() -> blockingCall());  // ✅ lazy (그래서 defer/fromCallable이 필요)

Backpressure

  • 소비자가 감당할 수 있는 양만 요청(request)하여 데이터 흐름을 조절해 메모리 오버플로우를 방지한다.
  • 실체는 Subscription.request(n) — 구독자가 N개를 더 달라고 신호를 보낸다.
  • 오버플로우 전략 : onBackpressureBuffer() / onBackpressureDrop() / onBackpressureLatest()
  • BaseSubscriber로 request(n)을 직접 제어할 수 있다.

Disposable & 리소스 관리

  • subscribe()가 반환하는 Disposable로 스트림을 취소할 수 있다.
  • dispose() — 상위(upstream)로 취소 신호를 전파한다.
  • isDisposed() — 구독 상태를 확인한다.
  • 무한 스트림(SSE, Kafka 리스너)에서 특히 중요
  • 취소 신호가 상위 생산자까지 전파되어 즉시 리소스 정리를 유발한다.
  • subscribe()를 하는 순간, 각 단계가 서로를 참조하는 연결 사슬이 만들어진다.
  • 여기서 cancel()을 호출하게 되면 상위 Subscription을 들고 있으므로 연쇄적으로 취소가 된다.

스레드 모델

  • 이벤트 루프 : WebFlux는 소수의 이벤트 루프 스레드(Netty, 보통 코어 수)로 모든 요청을 처리한다.
  • MVC는 요청 1개당 쓰레드 1개를 할당해 블로킹돼도 그 쓰레드만 멈추지만 WebFlux는 쓰레드 몇 개가 수천 요청을 번갈아 처리한다.
// ❌ 이벤트 루프에서 블로킹 금지(안티패턴 - 이벤트 루프 스레드를 점유해 전체 처리량이 붕괴)
User user = jdbcUserRepository.findById(id); // 블로킹 JDBC!
// ✅ 불가피하면 격리한다.
Mono.fromCallable(() -> jdbcUserRepository.findById(id))
    .subscribeOn(Schedulers.boundedElastic());                // boundedElastic() — 블로킹 IO 격리

블로킹 대상: JDBC, RestTemplate, Thread.sleep(), 동기 파일 IO 등.

Reactive Streams 신호 흐름

언제 WebFlux를 쓰는가(트레이드오프)

  • 리액티브는 자동으로 빠르지 않다. 적은 스레드로 고동시성을 버티는 게 본질이다.
  • 적합 : 고동시성 + I/O 바운드(게이트웨이/BFF), SSE/WebSocket 스트리밍, 논블로킹 스택 전체
  • 부적합 : 단순 CRUD, CPU 바운드, 팀 미숙 → MVC + 가상 스레드(Loom)가 더 단순한 경우가 많다.

📖 Code Philosophy

📖 Java

📖 Kotlin

📖 Coroutine

📖 Spring

📖 Spring Security

📖 Spring Batch

📖 Database

📖 MySQL

📖 Redis

📖 JPA

📖 QueryDsl

📖 MSA

📖 Kafka

📖 Apache Flink

  • [Apache Flink - Apache Flink Architecture]
  • [Apache Flink - Stream Processing]
  • [Apache Flink - Data Stream API & Window]
  • [Apache Flink - State Management]

📖 HTTP

📖 AWS

📖 Docker

📖 Kubernetes

📖 Github Actions

📖 Jenkins

📖 Nginx

📖 Monitoring

📖 Test(feat. Load Testing)

  • [Test - Load Testing Fundamentals]
  • [Test - Identifying Bottlenecks with Load Testing]
  • [Test - Resolving Bottlenecks and Improving Performance]

📖 Test(feat. Java)

📖 Spring AI

📖 gRPC

  • [gRPC - Writing .proto Files with Protocol Buffers]
  • [gRPC - Various Communication Patterns in gRPC]
  • [gRPC - gRPC Optimization Techniques and Advanced Features]

📖 TDD(Test-Driven-Development)

📖 PostgreSQL

  • [PostgreSQL - Docker만을 사용하는 경량화된 환경 구성 방법]
  • [PostgreSQL - PostgreSQL에서 제공하는 데이터 타입]
  • [PostgreSQL - PostgreSQI의 JSONB, 역인덱싱과 활용 방법]
  • [PostgreSQL - 데이터베이스 성능을 위한 최적화 패턴 및 전략]
  • [PostgreSQL - 트랜잭션과 ACID, Isolation 수준별 차이]
  • [PostgreSQL - Database Lock 교착상태와 읽기/쓰기 성능을 보장하는 MVCC 모델]
  • [PostgreSQL - pgvector와 벡터 저장, 유사도 검색 패턴 개념]
  • [PostgreSQL - 벡터 인덱스 최적화와 벡터 검색과 전문 검색 결합 패턴]
  • [PostgreSQL - PostgreSQL 플러그인]
  • [PostgreSQL - PostGIS - 공간 쿼리와 GIST 인덱스, 지리 타입과 공간 쿼리를 위한 타입과 기본 함수]
  • [PostgreSQL - pg_search - 검색 엔진 없이 텍스트 검색 구현과 주의사항]
  • [PostgreSQL - 단일 인스턴스 한계를 극복하는 분산 패턴과 스케줄링, 분산 환경 구축 방법]
  • [PostgreSQL - Citus - 분산 테이블과 분산 쿼리를 위한 Extension과 데이터 분산 처리]
  • [PostgreSQL - pg_cron - PostgreSQL로 구성하는 CronJob]
  • [PostgreSQL - 스케줄러 + 분산 처리를 동시에 도입하는 주기적 집계 쿼리 패턴]

📖 Workflow-Driven Techniques for Large-Scale Traffic Processing

  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka + Debezium을 활용한 CDC 패턴 설계]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Temporal을 활용한 워크플로우 패턴]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Docker와 경량 이미지를 활용한 환경 구축 방법]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka에서의 메시지 Delivery Guarantee]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - 실시간 동기화의 핵심 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - MySQL Binary Log 기반의 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Binary Log 기반의 CDC 구현 플랫폼 Debezium이란?]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture Best Practice와 주의사항]

📖 Reactive Programming

📖 ElasticSearch

  • [ElasticSearch - Search Fundamentals & How Search Works]
  • [ElasticSearch - Korean-Optimized Search]
  • [ElasticSearch - Mapping & Data Types]
  • [ElasticSearch - Common Search Features]
  • [ElasticSearch - Building a Product Search Engine with Elasticsearch]
  • [ElasticSearch - Deploying Elasticsearch with Elastic Cloud]
  • [ElasticSearch - Managing Documents]
  • [ElasticSearch - Analysis & Mapping]
  • [ElasticSearch - Search Fundamentals]
  • [ElasticSearch - Query Joins]
  • [ElasticSearch - Processing Search Results]
  • [ElasticSearch - Aggregations]
  • [ElasticSearch - Tips for Improving Search Results]
  • [ElasticSearch - Elasticsearch Clients]
  • [ElasticSearch - Understanding How Elasticsearch Works]
  • [ElasticSearch - Monitoring Elasticsearch]
  • [ElasticSearch - Elasticsearch Troubleshooting]

Clone this wiki locally