SubmissionPublisher
SubmissionPublisher (비동기 제출 게시자)
제출된(비-null) 항목을 닫힐 때까지 현재 구독자들에게 비동기적으로 발행하는 Flow.Publisher예요. 각 구독자는 드롭이나 예외가 없는 한 새로 제출된 항목을 같은 순서로 받아요.
본문
SubmissionPublisher를 쓰면 항목 생성기가 드롭 처리나 블로킹을 이용한 흐름 제어에 의존하면서 리액티브 스트림 규격에 맞는 Publisher처럼 동작할 수 있어요. 구독자에게 항목을 전달할 때는 생성자에서 준 Executor를 사용해요.
- 항목 생성기가 별도 쓰레드에서 돌고 구독자 수를 예측할 수 있다면
Executors.newFixedThreadPool(int)을 고려해요. - 그 외에는 기본값인
ForkJoinPool.commonPool()을 쓰는 게 일반적이에요.
버퍼링 덕에 생산자와 소비자는 짧은 시간 동안 서로 다른 속도로 동작할 수 있어요. 각 구독자는 독립적인 버퍼를 쓰며, 버퍼는 처음 사용될 때 만들어지고 필요에 따라 최대 용량까지 늘어나요.
버퍼가 가득 찼을 때의 동작은 submit과 offer가 서로 달라요.
submit은 리소스가 생길 때까지 블로킹해요. 가장 단순하지만 응답성이 떨어져요.offer는 항목을 (즉시 또는 제한 시간 내에) 드롭할 수 있고, 핸들러를 끼워 넣은 뒤 재시도할 기회를 줘요.- 구독자의 메서드가 예외를 던지면 그 구독은 취소돼요.
consume(Consumer) 메서드는 구독자가 모든 항목을 요청하고 처리하는 일반적인 경우를 간단히 지원해요.
class PeriodicPublisher<T> extends SubmissionPublisher<T> {
final ScheduledFuture<?> periodicTask;
final ScheduledExecutorService scheduler;
PeriodicPublisher(Executor executor, int maxBufferCapacity, Supplier<? extends T> supplier,
long period, TimeUnit unit) {
super(executor, maxBufferCapacity);
scheduler = new ScheduledThreadPoolExecutor(1);
periodicTask = scheduler.scheduleAtFixedRate(
() -> submit(supplier.get()), 0, period, unit);
}
public void close() {
periodicTask.cancel(false);
scheduler.shutdown();
super.close();
}
}
이 클래스는 항목을 생성하는 서브클래스의 기반으로도 편리하게 쓰여요. 코드와 시그니처는 원문 그대로 보존돼요.