Flow API(java.util.concurrent.Flow)는 Java 9에서 새롭게 도입된 기능입니다. 이 API를 활용하면 Publisher(발행자)와 Subscriber(구독자) 인터페이스가 서로 상호작용하며 원하는 작업을 수행하는 다양한 방식을 쉽게 이해할 수 있습니다.
Flow API의 핵심 구성 요소
Flow API는 반응형 스트림(Reactive Streams) 사양을 기반으로 하며, 다음 네 가지 핵심 인터페이스로 구성됩니다.
- Publisher: 데이터를 생성하고 발행하는 역할
- Subscriber: 발행된 데이터를 구독하여 소비하는 역할
- Subscription: Publisher와 Subscriber 간의 연결 관계를 관리
- Processor: Publisher와 Subscriber의 기능을 모두 수행하는 중간 처리자
아래 예제에서는 Publisher-Subscriber 인터페이스를 사용하여 Flow API를 직접 구현해 보겠습니다.
예제 코드
import java.util.concurrent.Flow.Publisher;
import java.util.concurrent.Flow.Subscriber;
import java.util.concurrent.Flow.Subscription;
public class FlowAPITest {
public static void main(String args[]) {
// Publisher 생성
Publisher<Integer> publisherSync = new Publisher<Integer>() {
@Override
public void subscribe(Subscriber<? super Integer> subscriber) {
for(int i = 0; i < 10; i++) {
System.out.println(Thread.currentThread().getName() + " | Publishing = " + i);
subscriber.onNext(i);
}
subscriber.onComplete();
}
};
// Subscriber 생성
Subscriber<Integer> subscriberSync = new Subscriber<Integer>() {
@Override
public void onSubscribe(Subscription subscription) {
}
@Override
public void onNext(Integer item) {
System.out.println(Thread.currentThread().getName() + " | Received = " + item);
try {
Thread.sleep(100);
} catch(InterruptedException e) {
e.printStackTrace();
}
}
@Override
public void onError(Throwable throwable) {
}
@Override
public void onComplete() {
}
};
publisherSync.subscribe(subscriberSync);
}
}실행 결과
main | Publishing = 0 main | Received = 0 main | Publishing = 1 main | Received = 1 main | Publishing = 2 main | Received = 2 main | Publishing = 3 main | Received = 3 main | Publishing = 4 main | Received = 4 main | Publishing = 5 main | Received = 5 main | Publishing = 6 main | Received = 6 main | Publishing = 7 main | Received = 7 main | Publishing = 8 main | Received = 8 main | Publishing = 9 main | Received = 9
코드 동작 방식 설명
위 예제의 동작 흐름은 다음과 같습니다.
- Publisher 구현: 익명 클래스 방식으로 Publisher를 생성하고, subscribe() 메서드 내에서 0부터 9까지의 정수를 순차적으로 발행합니다.
- Subscriber 구현: onNext() 메서드가 호출될 때마다 전달받은 데이터를 출력하고, Thread.sleep(100)을 통해 약간의 지연을 주어 데이터 처리 시간을 시뮬레이션합니다.
- 구독 시작: publisherSync.subscribe(subscriberSync) 호출로 두 객체가 연결되고, 데이터 발행이 시작됩니다.
실행 결과를 보면 모든 작업이 main 스레드에서 동기적으로 처리되고 있음을 확인할 수 있습니다. 즉, Publisher가 데이터를 발행하면 Subscriber가 해당 데이터를 받아 처리한 후에야 다음 데이터가 발행되는 순차적 구조입니다.
실무 환경에서는 SubmissionPublisher라는 기본 제공 구현체를 사용하거나, 백프레셔(backpressure)를 적절히 처리하여 비동기 방식으로 확장할 수 있습니다. 이를 통해 대량의 데이터 스트림을 효율적이고 안정적으로 처리할 수 있습니다.