본문 바로가기

Spring Framework/Spring WebFlux

[Spring WebFlux] 4편: Cold/Warm/Hot Sequence

1. 들어가며

같은 Flux를 두 번 구독했을 뿐인데, DB 쿼리가 두 번 나간 걸 본 적이 있으신가요?

Flux<User> users = userRepository.findAll(); // R2DBC 조회

users.subscribe(u -> log.info("A: {}", u));
users.subscribe(u -> log.info("B: {}", u));

 

로그를 열어보면 쿼리가 실행된 흔적이 두 번 찍혀 있습니다. 캐시를 해둔 것도 아닌데 왜 매번 새로 조회가 일어날까요? 반대로 어떤 경우엔 두 번째 구독자가 아예 아무 데이터도 못 받는 경우도 있습니다. 이 차이를 설명하는 개념이 바로 Cold Sequence와 Hot Sequence입니다.

 

2. Cold Sequence란

Cold Sequence는 구독(subscribe)할 때마다 시퀀스가 처음부터 새로 시작되는 것을 말합니다. 정확히는, 구독자마다 자기만의 독립적인 데이터 소스와 실행 흐름을 갖는다는 뜻입니다.

 

앞서 본 userRepository.findAll() 예시가 바로 Cold입니다. subscribe()가 호출되는 순간 비로소 쿼리가 실행되고, 두 번째 subscribe()는 완전히 별개의 쿼리 실행을 트리거합니다. 두 구독자는 서로의 존재를 전혀 모른 채, 각자 처음부터 끝까지 데이터를 받습니다.

Flux<Integer> cold = Flux.range(1, 3)
    .doOnSubscribe(s -> log.info("새로운 구독 시작!"));

cold.subscribe(i -> log.info("A: {}", i));
cold.subscribe(i -> log.info("B: {}", i));
새로운 구독 시작!
A: 1
A: 2
A: 3
새로운 구독 시작!
B: 1
B: 2
B: 3

doOnSubscribe가 두 번 찍힌다는 게 핵심입니다. Reactor의 Flux/Mono는 기본적으로 Cold입니다. HTTP 요청, DB 쿼리처럼 "구독자마다 독립적으로 처리되어야 자연스러운" 작업들이 대부분 Cold인 이유이기도 합니다.

 

3. Hot Sequence란

Hot Sequence는 반대로, 시퀀스가 구독자와 무관하게 이미 진행 중인 경우입니다. 구독자는 이 흐름에 뒤늦게 끼어드는 입장이고, 자신이 구독한 시점 이후에 발생한 이벤트만 받을 수 있습니다. 구독 이전에 이미 지나가 버린 이벤트는 다시 받을 방법이 없습니다.

 

마우스 클릭 이벤트, 센서 데이터 스트림, 주식 시세 브로드캐스트 같은 것들이 전형적인 Hot입니다. 이벤트는 구독자가 있든 없든 계속 발생합니다.

 

바꿔 말하면 Cold가 "구독자마다 실행을 복제"하는 것이라면, Hot은 "하나의 실행을 구독자들이 나눠 갖는" 것입니다. 원본 시퀀스는 단 하나뿐이고, 구독자들은 그 원본에 올라타는 입장입니다. 그래서 몇 명이 구독하든 업스트림에서 실제로 일어나는 작업(쿼리 실행, 이벤트 발생 등)은 딱 한 번뿐입니다 — Cold에서 구독자 수만큼 작업이 반복되는 것과 정반대입니다.

Sinks.Many<Integer> sink = Sinks.many().multicast().directBestEffort();
Flux<Integer> hot = sink.asFlux();

sink.tryEmitNext(1); // 아직 아무도 구독 안 함 → 유실
hot.subscribe(i -> log.info("A: {}", i));

sink.tryEmitNext(2);
sink.tryEmitNext(3);
A: 2
A: 3

1은 구독자가 없던 시점에 방출되어 그대로 사라졌습니다. 이게 Hot의 본질입니다 — 스트림은 구독 여부와 무관하게 자기 갈 길을 갑니다.

 

참고로 Sinks.many().multicast().onBackpressureBuffer()는 겉보기엔 비슷해 보이지만 다르게 동작합니다. 이건 첫 구독자가 나타나기 전에 emit된 값을 "warm up" 버퍼에 기억해뒀다가 첫 구독자에게 그대로 전달하는 스펙이라, 여기서 보여주려는 순수 Hot의 유실 동작과는 맞지 않습니다. 구독 전 이벤트가 정말로 사라지는 걸 보여주려면 directBestEffort()나 directAllOrNothing()을 써야 합니다.

 

4. Warm Sequence — Cold와 Hot 사이

실무에서 마주치는 시퀀스는 이 둘 중 하나로 딱 떨어지지 않는 경우가 많습니다. 그 사이에 있는 것이 Warm Sequence입니다.

  • 첫 구독자가 나타나기 전까지는 Cold처럼 대기합니다 (실행이 시작 안 됨)
  • 첫 구독이 들어오는 순간 실행이 시작되고, 그 시점부터는 Hot처럼 하나의 실행을 여러 구독자가 공유합니다
  • 구독자가 전부 떨어져 나가면 다시 Cold 상태로 리셋되는 경우도 있습니다 (연산자 옵션에 따라 다름)

즉 Warm은 "누가 최초로 트리거하는가"는 Cold의 성질을, "그 이후 실행을 공유하는가"는 Hot의 성질을 함께 갖습니다. Reactor에서 share()/publish()로 만들어지는 시퀀스 대부분이 사실 순수 Hot이 아니라 이 Warm에 해당합니다.

 

 

세 가지를 한 그림에 놓고 보면 차이가 더 분명합니다. Cold는 구독자마다 완전히 다른 실행 라인을 갖고, Warm은 첫 구독이 실행을 시작시키되 그 뒤로는 하나의 라인을 공유하며, Hot은 구독 여부와 상관없이 이미 흘러가고 있는 라인에 올라타는 구조입니다.

 

5. 코드로 Cold → Warm 전환 확인하기: share()

Reactor는 share() 또는 publish()로 Cold 시퀀스를 Warm하게 바꿀 수 있습니다. 다음은 Flux.interval(원래 Cold)에 share()를 적용해 여러 구독자가 하나의 실행 흐름을 공유하도록 만든 예시입니다.

Flux<Long> source = Flux.interval(Duration.ofMillis(100))
    .doOnNext(i -> log.info("emit: {}", i))
    .share(); // Warm으로 전환

source.subscribe(i -> log.info("A: {}", i));

Thread.sleep(300); // A만 구독한 채로 300ms 경과

source.subscribe(i -> log.info("B: {}", i)); // B는 뒤늦게 합류

 

share()가 없었다면 A와 B는 각자 0부터 시작하는 별개의 타이머를 갖게 됩니다. share()를 걸면 하나의 업스트림 실행을 공유하게 되고, 뒤늦게 합류한 B는 자신이 구독한 시점 이후의 값만 받습니다 — 이전에 지나간 0, 1, 2는 놓칩니다.

3번 섹션의 Sinks.Many 예시와 비교해보면 차이가 분명해집니다. 그쪽은 구독자 존재 여부와 무관하게 이미 방출이 진행되고 있었지만, 이쪽은 A라는 첫 구독이 있어야 비로소 실행이 시작됩니다. 그래서 이 예시는 순수 Hot이 아니라 Warm입니다.

 

6. 코드로 확인하기: share() vs cache()

cache()도 Cold를 전환시키는 연산자지만, 결과는 share()와 다릅니다.

Flux<Long> source = Flux.interval(Duration.ofMillis(100))
    .doOnNext(i -> log.info("emit: {}", i))
    .take(3)
    .cache(); // Warm으로 전환, 단 replay 방식

source.subscribe(i -> log.info("A: {}", i));

Thread.sleep(500); // A만 구독한 채로 500ms 경과 (0, 1, 2 모두 방출 완료)

source.subscribe(i -> log.info("B: {}", i)); // B는 뒤늦게 합류
emit: 0
A: 0
emit: 1
A: 1
emit: 2
A: 2
B: 0
B: 1
B: 2

첫 구독이 있어야 실행이 시작된다는 점(Cold적 성질)은 share()와 같습니다. 하지만 뒤늦게 합류한 B는 지나간 값을 놓치는 게 아니라, 이미 방출된 값을 처음부터 그대로 재생받습니다 — emit: 로그가 다시 찍히지 않는다는 점에서 실행 자체는 한 번뿐이었다는 것도 확인할 수 있습니다.

 

여기서 헷갈리기 쉬운 지점이 하나 있습니다. B가 결국 0, 1, 2를 처음부터 다 받으니, 이것도 그냥 Cold 아닌가 싶을 수 있습니다. 하지만 Cold/Warm/Hot을 가르는 기준은 "구독자가 어떤 값을 받는가"가 아니라 "업스트림 실행이 몇 번 일어나는가"입니다.

 

진짜 Cold였다면(cache() 없이 Flux.interval()을 그대로 두 번 구독했다면) emit: 로그가 A 구독으로 3번, B 구독으로 또 3번, 총 6번 찍혔을 것입니다. 하지만 위 예시에서는 emit: 로그가 딱 3번뿐입니다 — 업스트림 실행은 한 번만 일어났고, B는 그 결과를 재생받았을 뿐입니다. 값의 모양은 Cold와 비슷해 보여도, 실행이 반복되지 않는다는 점에서 Warm입니다.

 

 

정리하면 둘 다 Warm이지만 성격이 다릅니다.

  • share(): 실행은 공유하되, 늦게 온 구독자는 그 시점 이후 값만 받습니다 (유실 있음)
  • cache(): 실행은 공유하고, 늦게 온 구독자에게 과거 값까지 전부 재생해줍니다 (유실 없음)

 

replay()는 이 cache()의 동작을 더 세밀하게 제어할 수 있는 버전이라고 보면 됩니다 (버퍼 크기, TTL 등 옵션 지정 가능). 뒤에서 다시 다루겠지만, replay()는 cache()처럼 Cold를 Warm으로 바꾸는 용도뿐 아니라, 3번 섹션에서 본 것 같은 순수 Hot 스트림에 재생 버퍼를 붙이는 용도로도 씁니다 — 어느 쪽이든 "늦게 온 구독자에게 과거 값을 재생해준다"는 역할은 같습니다.

 

7. 왜 이 구분이 중요한가

이 구분을 모르고 넘어가면 실무에서 두 가지 방향의 함정에 빠지기 쉽습니다.

  • Cold를 여러 번 구독해서 리소스를 낭비하는 경우: 같은 Mono<Response>를 여러 소비자가 각자 구독하면, 외부 API 호출이나 DB 쿼리가 구독자 수만큼 반복됩니다. 결과를 공유해야 한다면 cache()나 share()로 Warm하게 만들어야 합니다.
  • Hot에 늦게 구독해서 이벤트를 유실하는 경우: 실시간 브로드캐스트 스트림에 뒤늦게 구독한 컴포넌트가 "왜 초기 데이터가 안 오지?"라고 의아해하는 상황입니다. 이건 버그가 아니라 Hot Sequence의 정상 동작입니다. 필요하다면 replay()로 최근 N개를 버퍼링해서 신규 구독자에게 재생해줘야 합니다.

즉 Cold/Warm/Hot은 단순한 용어 구분이 아니라, "이 스트림을 여러 번/여러 명이 구독해도 안전한가"를 판단하는 실질적인 기준입니다.

 

8. 다음 편 예고

share(), publish(), cache(), replay() 같은 연산자들이 실행을 공유시켜준다는 건 확인했지만, 정작 "구독자마다 처리 속도가 다르면 어떻게 되는가"라는 질문은 아직 남아 있습니다. Cold든 Warm이든 Hot이든, 생산자가 소비자보다 빠르게 데이터를 밀어내는 순간 문제가 생기는 건 마찬가지입니다. 다음 편에서는 Reactor가 이 문제를 request(n) 기반으로 어떻게 제어하는지, 그리고 감당 못 할 때 쓸 수 있는 배압 전략들을 다루겠습니다.

 

9. 마무리

Cold/Warm/Hot 구분은 사실 리액티브 프로그래밍만의 개념은 아닙니다 — "누가 데이터의 생명주기를 소유하는가"라는 오래된 질문의 리액티브 버전일 뿐입니다. 하지만 이 구분을 명확히 갖고 있어야 다음 편에서 다룰 배압 제어 — 즉 "속도 차이를 어떻게 다스릴 것인가" — 가 훨씬 자연스럽게 읽힙니다.