SubmissionPublisher

SubmissionPublisher (비동기 제출 게시자)

제출된(비-null) 항목을 닫힐 때까지 현재 구독자들에게 비동기적으로 발행하는 Flow.Publisher예요. 각 구독자는 드롭이나 예외가 없는 한 새로 제출된 항목을 같은 순서로 받아요.

출처: Java API Reference

본문

SubmissionPublisher를 쓰면 항목 생성기가 드롭 처리나 블로킹을 이용한 흐름 제어에 의존하면서 리액티브 스트림 규격에 맞는 Publisher처럼 동작할 수 있어요. 구독자에게 항목을 전달할 때는 생성자에서 준 Executor를 사용해요.

  • 항목 생성기가 별도 쓰레드에서 돌고 구독자 수를 예측할 수 있다면 Executors.newFixedThreadPool(int)을 고려해요.
  • 그 외에는 기본값인 ForkJoinPool.commonPool()을 쓰는 게 일반적이에요.

버퍼링 덕에 생산자와 소비자는 짧은 시간 동안 서로 다른 속도로 동작할 수 있어요. 각 구독자는 독립적인 버퍼를 쓰며, 버퍼는 처음 사용될 때 만들어지고 필요에 따라 최대 용량까지 늘어나요.

버퍼가 가득 찼을 때의 동작은 submitoffer가 서로 달라요.

  • 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();
    }
}

이 클래스는 항목을 생성하는 서브클래스의 기반으로도 편리하게 쓰여요. 코드와 시그니처는 원문 그대로 보존돼요.

더 알아보기