Java 9에서는 java.util.concurrent.Flow 패키지를 통해 리액티브 스트림즈(Reactive Streams)가 공식적으로 도입되었습니다. 이는 상호 운용 가능한 발행-구독(Publish-Subscribe) 프레임워크를 지원하며, 비동기 경계(다른 스레드 또는 스레드 풀로 요소를 전달하는 과정)를 넘어 비동기 데이터 스트림을 처리할 수 있게 해줍니다.
리액티브 스트림즈의 가장 큰 특징은 역압력(Backpressure) 메커니즘입니다. 수신 측이 임의의 양의 데이터를 강제로 버퍼링하지 않고 처리 가능한 만큼만 요청하기 때문에, 버퍼 오버플로우(Buffer Overflow)가 발생하지 않습니다.
Flow API의 4가지 핵심 인터페이스
Flow API는 서로 밀접하게 연관된 네 가지 핵심 인터페이스로 구성됩니다.
- Publisher – 데이터를 발행하는 생산자
- Subscriber – 데이터를 소비하는 소비자
- Subscription – Publisher와 Subscriber 사이의 연결 관계를 관리
- Processor – Subscriber이자 Publisher 역할을 동시에 수행하는 중개자
인터페이스 구문(Syntax)
@FunctionalInterface
public static interface Publisher<T> {
public void subscribe(Subscriber<? super T> subscriber);
}
public static interface Subscriber<T> {
public void onSubscribe(Subscription subscription);
public void onNext(T item);
public void onError(Throwable throwable);
public void onComplete();
}
public static interface Subscription {
public void request(long n);
public void cancel();
}
public static interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}각 인터페이스의 주요 메서드
1. Flow.Publisher
subscribe() 메서드 하나만을 가지며, 해당 Publisher에 Subscriber를 등록하는 역할을 합니다.
2. Flow.Subscriber
네 가지 콜백 메서드를 제공합니다.
onSubscribe()– 구독이 시작될 때 호출되며 Subscription 객체를 전달받습니다.onNext()– 새로운 데이터 항목이 도착할 때마다 호출됩니다.onError()– 처리 중 오류가 발생했을 때 호출됩니다.onComplete()– 모든 데이터 발행이 정상적으로 완료되었을 때 호출됩니다.
3. Flow.Subscription
역압력 제어의 핵심으로 두 가지 메서드를 가집니다.
request(long n)– Subscriber가 처리할 수 있는 만큼의 데이터 개수(n개)를 요청합니다.cancel()– 구독을 취소하여 더 이상 데이터를 받지 않도록 합니다.
4. Flow.Processor
Flow.Publisher와 Flow.Subscriber 인터페이스를 모두 확장하므로, 두 인터페이스의 모든 메서드를 구현해야 합니다. 이를 통해 데이터를 변환·중계하는 파이프라인 단계로 활용할 수 있습니다.
마무리
이처럼 Java 9의 Flow.Publisher, Flow.Subscriber, Flow.Subscription, Flow.Processor 네 가지 인터페이스는 리액티브 스트림즈 명세(Reactive Streams Specification)를 기반으로 설계되었습니다. 이 표준 API 덕분에 RxJava, Project Reactor 같은 외부 라이브러리와 JDK 간의 상호 운용성이 크게 향상되었으며, 비동기 환경에서도 안전하고 효율적인 데이터 스트림 처리가 가능해졌습니다.