Reactor 구성 요소
public class HelloReactorExample {
public static void main(String[] args) throws InterruptedException {
Flux // 여러 건의 데이터를 처리함을 의미
.just("Hello", "Reactor") // 원본 데이터 소스로부터 데이터를 emit 하는 publisher 역할
.map(message -> message.toUpperCase()) // map을 이용하여 대문자로 변경
.publisOn(Schedulers.parallel()) // 쓰레드 관리자 역할을 하는 Scheduler를 지정
.subscribe(System.out::println, // Publisher가 emit한 데이터를 전달 받아서 처리
error -> System.out.println(error.getMessage()), // 에러가 발생할 경우, 해당 에러를 전달 받아서 처리하는 역할
() -> System.out.println("# onComplete")); // Reactor Sequence가 종료된 후 어떤 후처리를 하는 역할
Thread.sleep(100L);
}
}
스케줄러
쓰레드를 관리하는 관리자의 역할
public class Schedulers {
public static void main(String[] args) throws InterruptedException {
Flux
.range(1, 10)
.doOnSubscribe(subscription -> log.info("# doOnSubscribe")) // 구독 직후 싱행되는 쓰레드와 동일한 쓰레드에서 실행
.subscribeOn(Schedulers.boundedElastic()) // 구독 직후에 실행되는 쓰레드가 main 쓰레드에서 Scheduler로 지정한 쓰레드로 바뀜
.filter(n -> n % 2 == 0)
.map(n -> n * 2)
.subscribe(data -> log.info("# onNext: {}", data));
Thread.sleep(100L);
}
}
public class SchedulersExample03 {
public static void main(String[] args) throws InterruptedException {
Flux
.range(1, 10)
.subscribeOn(Schedulers.boundedElastic())
.doOnSubscribe(subscription -> log.info("# doOnSubscribe"))
.publishOn(Schedulers.parallel()) // publishOn()에서 Scheduler로 지정한 쓰레드로 변경
.filter(n -> n % 2 == 0)
.doOnNext(data -> log.info("# filter doOnNext")) // 어느 쓰레드에서 실행이 되는지 확인하기 위한 용도
.publishOn(Schedulers.parallel()) // publishOn()에서 Scheduler로 지정한 쓰레드로 변경
.map(n -> n * 2)
.doOnNext(data -> log.info("# map doOnNext")) // 어느 쓰레드에서 실행이 되는지 확인하기 위한 용도
.subscribe(data -> log.info("# onNext: {}", data));
Thread.sleep(100L);
}
}
subscribeOn() : 데이터 소스에서 데이터를 emit하는 원본 Publisher의 실행 쓰레드를 지정하는 역할
publishOn() : 전달 받은 데이터를 가공 처리하는 Operator 앞에 추가해서 실행 쓰레드를 별도로 추가하는 역할