Computer >> 컴퓨터 >  >> 프로그래밍 >> Java

Java 9 Flow API 완벽 가이드: Publisher-Subscriber로 반응형 스트림 구현하기

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

코드 동작 방식 설명

위 예제의 동작 흐름은 다음과 같습니다.

  1. Publisher 구현: 익명 클래스 방식으로 Publisher를 생성하고, subscribe() 메서드 내에서 0부터 9까지의 정수를 순차적으로 발행합니다.
  2. Subscriber 구현: onNext() 메서드가 호출될 때마다 전달받은 데이터를 출력하고, Thread.sleep(100)을 통해 약간의 지연을 주어 데이터 처리 시간을 시뮬레이션합니다.
  3. 구독 시작: publisherSync.subscribe(subscriberSync) 호출로 두 객체가 연결되고, 데이터 발행이 시작됩니다.

실행 결과를 보면 모든 작업이 main 스레드에서 동기적으로 처리되고 있음을 확인할 수 있습니다. 즉, Publisher가 데이터를 발행하면 Subscriber가 해당 데이터를 받아 처리한 후에야 다음 데이터가 발행되는 순차적 구조입니다.

실무 환경에서는 SubmissionPublisher라는 기본 제공 구현체를 사용하거나, 백프레셔(backpressure)를 적절히 처리하여 비동기 방식으로 확장할 수 있습니다. 이를 통해 대량의 데이터 스트림을 효율적이고 안정적으로 처리할 수 있습니다.