Java 9는 Reactive Streams(리액티브 스트림)를 지원하기 위해 몇 가지 인터페이스를 새롭게 도입했습니다. 바로 Publisher, Subscriber, Subscription 인터페이스와, Publisher 인터페이스를 구현하는 SubmissionPublisher 클래스입니다. 각 인터페이스는 Reactive Streams의 원칙에 따라 서로 다른 역할을 담당합니다.
Subscriber 인터페이스는 퍼블리셔(Publisher)가 발행하는 데이터를 구독할 때 사용합니다. 이를 활용하려면 Subscriber 인터페이스를 직접 구현하고, 추상 메서드들에 대한 구현체를 제공해야 합니다.
Flow.Subscriber 인터페이스의 주요 메서드
- onComplete(): Publisher 객체가 자신의 역할을 모두 완료했을 때 호출됩니다.
- onError(): Publisher에서 오류가 발생하여 Subscriber에게 알림이 전달될 때 호출됩니다.
- onNext(): Publisher가 새로운 데이터를 발행하여 모든 Subscriber에게 알릴 때마다 호출됩니다.
- onSubscribe(): Publisher가 Subscriber를 등록(구독 추가)했을 때 호출됩니다.
예제 코드
아래 예제에서는 Flow.Subscriber<Integer>를 구현한 내부 클래스를 정의하고, SubmissionPublisher를 통해 1부터 10까지의 숫자를 발행하여 처리하는 과정을 보여줍니다.
import java.util.concurrent.Flow;
import java.util.concurrent.SubmissionPublisher;
import java.util.stream.IntStream;
public class SubscriberImplTest {
public static class Subscriber implements Flow.Subscriber<Integer> {
private Flow.Subscription subscription;
private boolean isDone;
@Override
public void onSubscribe(Flow.Subscription subscription) {
System.out.println("Subscribed");
this.subscription = subscription;
this.subscription.request(1);
}
@Override
public void onNext(Integer item) {
System.out.println("Processing " + item);
this.subscription.request(1);
}
@Override
public void onError(Throwable throwable) {
throwable.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("Processing done");
isDone = true;
}
}
public static void main(String args[]) throws InterruptedException {
SubmissionPublisher<Integer> publisher = new SubmissionPublisher<>();
Subscriber subscriber = new Subscriber();
publisher.subscribe(subscriber);
IntStream intData = IntStream.rangeClosed(1, 10);
intData.forEach(publisher::submit);
publisher.close();
while(!subscriber.isDone) {
Thread.sleep(10);
}
System.out.println("Done");
}
}실행 결과
Subscribed Processing 1 Processing 2 Processing 3 Processing 4 Processing 5 Processing 6 Processing 7 Processing 8 Processing 9 Processing 10 Processing done Done
코드 설명
위 코드의 동작 흐름은 다음과 같습니다.
- 구독 등록: publisher.subscribe(subscriber)가 호출되면 Subscriber의 onSubscribe() 메서드가 실행되고, 이때 subscription.request(1)을 통해 처음 한 개의 데이터를 요청합니다.
- 데이터 처리: Publisher가 submit()으로 데이터를 발행하면 onNext()가 호출되어 각 항목을 처리한 뒤, 다시 request(1)로 다음 데이터를 요청하는 백프레셔(backpressure) 방식으로 동작합니다.
- 완료 처리: publisher.close()가 호출되면 스트림이 종료되고, onComplete() 메서드가 실행되어 모든 처리가 끝났음을 알립니다.
이처럼 Java 9의 Flow API를 활용하면 비동기 데이터 스트림을 안정적으로 처리할 수 있으며, 요청 기반(request-based)의 백프레셔 제어를 통해 생산자와 소비자 간의 데이터 흐름을 효율적으로 관리할 수 있습니다.