2026. 2. 5. 01:22ㆍspring/Reactor
Programmatically creating a sequence
이 섹션에서는 Flux나 Mono를 생성할 때,
관련된 이벤트(onNext, onError, onComplete)를 프로그래밍적으로 정의하여 만드는 방법을 소개한다.
이 메서드들은 이벤트를 발생시키는 API 를 제공하는데 sink 라고 부른다.
sink 는 몇 가지 종류가 있는데 아래에서 살펴볼 것이다.
1. 동기식 생성: generate
Flux를 프로그래밍적으로 생성하는 가장 단순한 방법은 generate 메서드를 사용하는 것인데,
이 메서드는 generator 함수를 인자로 받는다.
이것은 동기적이고 요소를 하나씩 내보내는 용도이다.
sink 는 SynchronousSink 타입이고, next() 메서드는 콜백 호출 시 최대 한 번만 호출할 수 있습니다
추가로 error 나 complete()를 호출할 수도 있지만, 이는 선택 사항이다.
가장 유용한 메서드는 sink를 사용할 때 참조할 수 있는 상태(state)를 유지할 수 있게 해주는 메서드다.
이를 통해 다음에 어떤 데이터를 내보낼지 결정할 수 있다.
generator 함수는 BiFunction<S, SynchronousSink<T, S> 가 되고, 여기서 <S>는 상태 객체(state)의 타입이다.
초기 상태를 Supplier<S> 로 제공해야 하며, generator 함수는 각 단계마다 새로운 상태를 반환한다.
예를 들어 int로 상태를 사용할 수 있다.
// state-based example
Flux<String> flux = Flux.generate(
// 초기 상태로 0을 제공한다.
() -> 0,
(state, sink) -> {
// 상태를 보고 방출할 데이터를 정한다.(3의 곱셈표에서 행이 데이터다.)
sink.next("3 x " + state + " = " + 3*state);
// 상태를 보고 언제 멈출지도 정한다.
if (state == 10) sink.complete();
// 다음 호출에 사용할 새로운 상태를 반환한다.
return state + 1;
});
=> 출력
3 x 0 = 0
3 x 1 = 3
3 x 2 = 6
3 x 3 = 9
3 x 4 = 12
3 x 5 = 15
3 x 6 = 18
3 x 7 = 21
3 x 8 = 24
3 x 9 = 27
3 x 10 = 30
mutable(변경 가능한) 객체로 S 를 사용해도 된다.
위 예제는 AtomicLong 을 상태로 다시 작성할 수 있다.
// Mutable state variant
Flux<String> flux = Flux.generate(
// 이번엔 상태로 mutable 객체를 만든다.
AtomicLong::new,
(state, sink) -> {
// 여기서 상태를 변경한다.
long i = state.getAndIncrement();
sink.next("3 x " + i + " = " + 3*i);
if (i == 10) sink.complete();
// 같은 인스턴스를 새로운 상태로 반환한다.
return state;
});
상태 객체가 어떤 자원(resource)을 정리해야 한다면, generate(Supplier<S>, BiFunction, Consumer<S>) 변형을 사용해서 마지막 상태 인스턴스를 정리(clean up)할 수 있다.
* 여기서 자원은 상태를 쓰고 난 뒤 메모리나 파일, 네트워크 같은 자원을 의미
아래 예제는 Consumer 를 포함한 generate 메서드다.
Flux<String> flux = Flux.generate(
AtomicLong::new,
// 상태에 mutable 객체를 사용한다.
(state, sink) -> {
// 상태를 변경한다.
long i = state.getAndIncrement();
sink.next("3 x " + i + " = " + 3*i);
if (i == 10) sink.complete();
// 같은 인스턴스를 새로운 상태로 반환한다.
return state;
// 마지막 상태 값인 11이 찍힌다.
}, (state) -> System.out.println("state: " + state));
=> 출력
3 x 0 = 0
3 x 1 = 3
3 x 2 = 6
3 x 3 = 9
3 x 4 = 12
3 x 5 = 15
3 x 6 = 18
3 x 7 = 21
3 x 8 = 24
3 x 9 = 27
3 x 10 = 30
state: 11
상태에 데이터베이스 커넥션이나 기타 프로세스 종료 시 처리해야 할 자원이 포함되어 있는 경우,
Consumer 람다는 커넥션을 닫거나 프로세스가 끝날 때 수행되어야 할 다른 작업들을 처리할 수 있다.
2. 비동기, 멀티 스레드: create
Flux 를 프로그래밍적으로 생성하는 진보된 방법은 create 다.
이 메서드는 각 단계마다 여러 요소를 방출하는데 적합하다.(멀티 스레드 환경에서도 안전)
create는 next, error, complete 메서드를 가진 FluxSink를 노출한다.
generate와는 달리 상태(state)를 기반으로 하는 변형은 제공하지 않는다.
대신, 콜백 안에서 멀티 스레드로 이벤트를 발생시키는 것은 가능하다.
create는 리스너 기반의 비동기 API처럼,
기존에 존재하는 API를 리액티브 세계와 연결하는 데 매우 유용하다.
*주의*
create 는 코드를 병렬화 하거나 비동기로 만드는 게 아니다.
만약 create 람다 안에서 블로킹 작업을 수행하면, 데드락 같은 상황에 빠질 수 있다.subscribeOn 을 사용하더라도 주의할 점이 있는데,예를 들어 sink.next(t) 를 무한 루프로 호출하는 것처럼 오래 블로킹 되는 create 람다는 파이프라인을 멈출 수 있다. *(1)
요청(request)이 실행되어야 할 스레드를 해당 루프가 계속 점유해 버리기 때문에, 요청 자체가 수행되지 않게 되기 때문이다.
이런 상황에서는 subscribeOn(Scheduler, false) 을 사용해야 한다.
여기서 requestOnSeparateThread = false로 설정하면,
create는 Scheduler의 스레드에서 실행되면서도
요청(request)은 원래의 스레드에서 수행되어 데이터 흐름이 계속 유지된다.
// *(1)
Flux<Integer> flux = Flux.create(sink -> {
// 오래 블로킹되는 create 람다
while (true) {
sink.next(1); // 무한히 데이터 밀어넣기
}
})
.subscribeOn(Schedulers.single());
flux
.doOnNext(v -> System.out.println("받음: " + v))
.subscribe();
리스너 API 를 사용한다고 가정해보자.
이 API 는 데이터를 청크 단위로 처리하고 아래 예제의 MyEventListener 인터페이스에 보이는 것처럼 두 개의 이벤트를 가진다.
1. 청크 데이터 준비
2. 처리 완료
interface MyEventListener<T> {
void onDataChunk(List<T> chunk);
void processComplete();
}
create 를 사용해 Flux<T> 에 이 인터페이스를 연결 할 수 있다.
Flux<String> bridge = Flux.create(sink -> {
// 이 이벤트 등록은 구독시 한 번 발생한다.
myEventProcessor.register(
new MyEventListener<String>() {
// myEventProcessor 에서 등록된 이벤트 리스터의 onDataChunk 이벤트가 발생하면
// 인자로 넘어온 chunk 의 원소들이 Flux 의 원소가 된다.
public void onDataChunk(List<String> chunk) {
for(String s : chunk) {
sink.next(s);
}
}
// myEventProcessor 에서 등록된 이벤트 리스터의 processComplete 이벤트가 발생하면
// 처리가 종료된다.
public void processComplete() {
sink.complete();
}
});
});
또한 create 는 비동기 API 들을 연결할 수 있고, 백프레셔를 관리하기 때문에
OverFlowStrategy 를 지정하여 백프레셔 상황에서 동작을 조정할 수 있다.
- IGNORE : 다운스트림의 백프레셔 요청을 완전히 무시한다.(방출된 데이터가 다운스트림으로 가지 못하고 큐에 쌓여 에러 발생 가능)
- ERROR : downstream이 처리 속도를 따라오지 못할 때 IllegalStateException을 발생시킨다.(데이터 유실 방지)
- DROP : downstream이 받을 준비가 되어 있지 않으면 들어오는 데이터를 버린다.
- LATEST : downstream이 upstream에서 온 신호 중 가장 최신 값만 받도록 한다.(최신 데이터 중시)
- BUFFER : 기본값으로, downstream이 처리하지 못하는 모든 신호를 버퍼에 저장한다. 이 방식은 무제한 버퍼링을 하기 때문에 OutOfMemoryError가 발생할 수 있다.
Mono에도 create 생성자(generator) 가 있다. Mono.create에서 사용하는 MonoSink는 여러 번의 방출을 허용하지 않는다.
첫 번째 신호 이후에 발생하는 모든 신호는 무시한다.
3. 비동기, 싱글 스레드: push
push 는 generate 와 create 의 중간 형태로,
단일 생산자의 이벤트 처리에 적합하다.
create 와 유사하게 비동기를 지원하며,
create 가 지원하는 오버플로우 전략을 사용해 백프레셔를 관리할 수 있다.
하지만 한 번에 하나의 생산자 스레드만 next, complete, error를 호출할 수 있다.
Flux<String> bridge = Flux.push(sink -> {
myEventProcessor.register(
// SingleThreadEventListener API 를 연결한다
new SingleThreadEventListener<String>() {
public void onDataChunk(List<String> chunk) {
for(String s : chunk) {
// 단일 리스너 스레드에서 sink 를 사용해 이벤트를 푸시한다.
sink.next(s);
}
}
public void processComplete() {
// 동일한 리스너 스레드에서 complete 이벤트를 발생시킨다.
sink.complete();
}
public void processError(Throwable e) {
// 동일한 리스너 스레드에서 error 이벤트를 발생시킨다.
sink.error(e);
}
});
});
3.1 push/pull 하이브리드 모델
create 같은 대부분의 리액터 연산자들은 hybrid push/pull 모델이다.
어떤 의미냐면, 대부분의 처리는 비동기일지라도(push 방식 암시),
작게나마 pull 요소가 존재한다. 바로 request 다
소비자가 먼저 요청하지 않으면 데이터를 방출하지 않는다는 의미에서, 소비자가 데이터를 pull 한다고 볼 수 있다.
소스는 데이터가 준비되는 즉시 소비자에게 데이터를 push하지만,
그 양은 소비자가 요청한 개수를 넘지 않는 범위 안에서만 전달한다.
push(), create() 둘 다 onRequest 소비자를 설정할 수 있는데,요청 양을 관리하고, 아직 처리되지 않은 요청이 있을 때만 sink를 통해 데이터가 푸시되도록 보장하기 위함이다.
Flux<String> bridge = Flux.create(sink -> {
myMessageProcessor.register(
new MyMessageListener<String>() {
// request 시점에는 없었던
// (3) 나중에 비동기로 도착한 메세지들도 sink 에 push 한다.
public void onMessage(List<String> messages) {
for(String s : messages) {
sink.next(s);
}
}
});
sink.onRequest(n -> {
// * 소비자가 request(n) 요청을 보낸 것임
// (1) 요청이 들어오면 메세지를 조회한다.
List<String> messages = myMessageProcessor.getHistory(n);
for(String s : messages) {
// (2) 메세지가 있다면 sink 에 push 한다.
sink.next(s);
}
});
});
3.2 push(), create() 후 cleaning up
onDipose, onCancel 두 콜백은 취소나 종료 시에 필요한 정리 작업을 수행한다.
onDispose는
Flux가 정상 완료되었을 때, 에러로 종료되었을 때, 취소되었을 때
정리 작업을 수행하는데 사용될 수 있다.
onCancel은
취소(cancellation)인 경우에만 실행되는 작업을 수행할 수 있으며,
그 실행 시점은 onDispose에 의한 정리 작업이 실행되기 이전이다.
Flux<String> bridge = Flux.create(sink -> {
sink.onRequest(n -> channel.poll(n))
// 취소 신호에만 호출되고 첫 번째로 실행된다.
.onCancel(() -> channel.cancel())
// 완료, 에러, 취소 신호에 호출된다.
.onDispose(() -> channel.close())
});
4. Handle
handle 메서드는 조금 다르다.
이것은 인스턴스 메서드이며,
Flux나 Mono 같은 이미 존재하는 소스에 체인으로 붙여서 사용하는 연산자다.
(map, filter 같은 일반 연산자들처럼)
handle은 generate와 비슷한데,
그 이유는 SynchronousSink를 사용하고,
한 번의 호출에서 하나의 값만 방출할 수 있기 때문이다.
하지만 handle은
각 소스 요소 하나당 임의의 값을 생성할 수 있고,
원하면 아예 방출하지 않고 건너뛸 수도 있다.
이런 특성 때문에 handle은 map + filter를 합쳐 놓은 연산자처럼 사용할 수 있다.
아래는 handle 메서드의 시그니쳐다
Flux<R> handle(BiConsumer<T, SynchronousSink<R>>);
예를 들어보자. Reactive Streams 명세에서는 시퀀스 안에 null 값이 들어오는 것을 허용하지 않는다.
그런데 map 연산을 하고 싶은데,
이미 존재하는 메서드를 map 함수로 쓰고 싶고,
그 메서드가 null을 반환할 수 있다면 어떻게 해야 할까?
예를 들어, 아래와 같은 메서드는 정수를 소스로 하는 스트림에 안전하게 적용할 수 있다.
public String alphabet(int letterNumber) {
if (letterNumber < 1 || letterNumber > 26) {
return null;
}
int letterIndexAscii = 'A' + letterNumber - 1;
return "" + (char) letterIndexAscii;
}
handle 을 사용해서 null 을 모두 제거할 수 있다.
Flux<String> alphabet = Flux.just(-1, 30, 13, 9, 20)
.handle((i, sink) -> {
// String 으로 변환
String letter = alphabet(i);
// 널 체크
if (letter != null)
// 필터
sink.next(letter);
});
alphabet.subscribe(System.out::println);
=> 출력
M
I
T
filter와 map 을 사용하면 되는거 아니냐 의문이 들 수 있다.
source
.filter(v -> isValid(v)) // null 여부를 미리 알 수 있음
.map(v -> convert(v)); // 절대 null 안 나옴
이런 경우는 괜찮다. 명세 위반도 없고 읽기도 쉽다.
하지만 아래와 같이 반환하기 전까진 null 이 나올지 안 나올지 모르는 경우가 있다.
source
.filter(v -> ???) // 여기서 null 여부 판단 불가
.map(v -> convert(v)); // 여기서 null 가능
물론 억지로 가능은 하다.
source
.filter(v -> convert(v) != null)
.map(v -> convert(v));
하지만 의도가 불명확하고 convert 메서드도 두 번 호출된다.
일관성도 고려해야 하고 비용 문제도 있을 수 있다.
'spring > Reactor' 카테고리의 다른 글
| Reactor(8) Sinks (0) | 2026.02.10 |
|---|---|
| Reactor(6) Thread, Scheduler, publishOn, subscribeOn (0) | 2026.02.07 |
| Reactor(3) Flux, Mono (0) | 2026.02.01 |
| Reactor(2) 리액티브 프로그래밍 소개 (1) | 2026.02.01 |
| Reactor(1) 시작하기 (0) | 2026.01.28 |