Flow — 리액티브 스트림 기반 흐름 제어 인터페이스

Flow — 리액티브 스트림 기반 흐름 제어 인터페이스

Flow 클래스는 흐름 제어(flow-controlled) 컴포넌트를 구축하기 위한 상호 관련된 인터페이스와 정적 메서드를 제공해요. Publisher가 항목을 생산하고, 하나 이상의 Subscriber가 그것을 소비하며, 각각을 Subscription이 관리해요. 이 인터페이스들은 reactive-streams 명세에 대응해요.

출처: Java API Reference

본문

개념 이해하기

Flow는 흐름 제어된 컴포넌트를 위한 상호 관련된 인터페이스와 static 메서드를 정의해요. Publisher가 항목을 생산하고 Subscriber가 소비하며 Subscription이 관리해요.

public final class Flow
extends Object

이 인터페이스들은 reactive-streams 명세에 대응하며, 동시성·분산 비동기 설정 모두에 적용돼요. 모든 일곱 메서드void "일방향(one-way)" 메시지 스타일로 정의돼요. 통신은 Flow.Subscription.request(long)이라는 단순한 흐름 제어에 의존해, "push" 기반 시스템에서 생기는 자원 관리 문제를 피해요.

주요 구성 요소

  • Flow.Publisher<T> — Subscriber에게 항목을 게시하는 생산자.
  • Flow.Subscriber<T> — 항목을 받고 처리하는 소비자. onSubscribe, onNext, onError, onComplete를 구현해요.
  • Flow.Subscription — Publisher와 Subscriber를 연결하는 1:1 제어/요청 채널. request(long)cancel()을 제공해요.
  • Flow.Processor<T,R> — Publisher이자 Subscriber인 컴포넌트.

간단한 Publisher 예

class OneShotPublisher implements Publisher<Boolean> {
    private final ExecutorService executor = ForkJoinPool.commonPool(); // daemon-based
    private boolean subscribed; // true after first subscribe
    public synchronized void subscribe(Subscriber<? super Boolean> subscriber) {
        if (subscribed)
            subscriber.onError(new IllegalStateException()); // only one allowed
        else {
            subscribed = true;
            subscriber.onSubscribe(new OneShotSubscription(subscriber, executor));
        }
    }
    static class OneShotSubscription implements Subscription {
        private final Subscriber<? super Boolean> subscriber;
        private final ExecutorService executor;
        private Future<?> future; // to allow cancellation
        private boolean completed;
        OneShotSubscription(Subscriber<? super Boolean> subscriber, ExecutorService executor) {
            this.subscriber = subscriber;
            this.executor = executor;
        }
        public synchronized void request(long n) {
            if (!completed) {
                completed = true;
                if (n <= 0) {
                    IllegalArgumentException ex = new IllegalArgumentException();
                    executor.execute(() -> subscriber.onError(ex));
                } else {
                    future = executor.submit(() -> {
                        subscriber.onNext(Boolean.TRUE);
                        subscriber.onComplete();
                    });
                }
            }
        }
        public synchronized void cancel() {
            completed = true;
            if (future != null) future.cancel(false);
        }
    }
}

Subscriber 예 — 요청량 조절

SampleSubscriber는 버퍼 크기에 따라 항목을 요청하고, 절반 소비 시점에 재요청해요. 주어진 Subscription에 대한 Subscriber 메서드 호출은 엄격히 순서화되므로 락이나 volatile이 필요 없어요(단 하나의 Subscriber가 여러 Subscription을 유지하지 않는 한).

class SampleSubscriber<T> implements Subscriber<T> {
    final Consumer<? super T> consumer;
    Subscription subscription;
    final long bufferSize;
    long count;
    SampleSubscriber(long bufferSize, Consumer<? super T> consumer) {
        this.bufferSize = bufferSize;
        this.consumer = consumer;
    }
    public void onSubscribe(Subscription subscription) {
        long initialRequestSize = bufferSize;
        count = bufferSize - bufferSize / 2; // re-request when half consumed
        (this.subscription = subscription).request(initialRequestSize);
    }
    public void onNext(T item) {
        if (--count <= 0)
            subscription.request(count = bufferSize - bufferSize / 2);
        consumer.accept(item);
    }
    public void onError(Throwable ex) { ex.printStackTrace(); }
    public void onComplete() {}
}

흐름 제어가 필요 없다면 사실상 무경계로 요청할 수 있어요:

class UnboundedSubscriber<T> implements Subscriber<T> {
    public void onSubscribe(Subscription subscription) {
        subscription.request(Long.MAX_VALUE); // effectively unbounded
    }
    public void onNext(T item) { use(item); }
    public void onError(Throwable ex) { ex.printStackTrace(); }
    public void onComplete() {}
    void use(T item) { ... }
}

정적 메서드

public static int defaultBufferSize() — 다른 제약이 없을 때 Publisher나 Subscriber 버퍼링에 쓸 기본 값을 반환해요. 현재 값은 256이에요. 예상 속도·자원·용법에 따라 요청 크기와 용량을 고를 때 유용한 출발점이 돼요.

더 알아보기 (Learn more)