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

Java 9 Flow API로 리액티브 스트림(Reactive Streams) 구현하는 방법

Flow API란 무엇인가?

Flow API는 Java 9부터 도입된 리액티브 스트림(Reactive Streams) 명세의 공식 지원 기능입니다. 이 API는 이터레이터(Iterator) 패턴과 옵저버(Observer) 패턴을 결합한 형태로, 데이터 항목을 비동기적으로 처리하면서도 백프레셔(backpressure)를 효과적으로 관리할 수 있게 해줍니다.

다만 Flow API는 RxJava처럼 최종 사용자를 위한 완성형 라이브러리가 아니라, 서로 다른 리액티브 라이브러리 간의 상호 운용성(interoperability)을 보장하기 위한 표준 명세라는 점에 유의해야 합니다.

Flow API의 4가지 핵심 인터페이스

  • Publisher(발행자): 등록된 구독자들에게 데이터 항목 스트림을 발행합니다.
  • Subscriber(구독자): Publisher를 구독하여 콜백을 통해 데이터를 수신합니다.
  • Subscription(구독): Publisher와 Subscriber를 연결하는 링크 역할을 하며, 데이터 요청(request)과 구독 취소(cancel)를 관리합니다.
  • Processor(프로세서): Publisher와 Subscriber 사이에 위치하여 한 스트림을 다른 스트림으로 변환합니다.

예제 코드

아래 예제에서는 하나의 데이터 객체를 요청하고, 이를 출력한 뒤 다시 하나를 추가로 요청하는 기본적인 구독자(MySubscriber)를 작성합니다. 그리고 Java에서 기본 제공하는 Publisher 구현체인 SubmissionPublisher를 활용해 전체 발행-구독 흐름을 완성할 수 있습니다.

import java.util.concurrent.Flow;
import java.util.List;
import java.util.concurrent.SubmissionPublisher;

class MySubscriber<T> implements Flow.Subscriber<T> {
    private Flow.Subscription subscription;

    @Override
    public void onSubscribe(Flow.Subscription subscription) {
        this.subscription = subscription;
        this.subscription.request(1);
    }

    @Override
    public void onNext(T item) {
        System.out.println(item);
        subscription.request(1);
    }

    @Override
    public void onError(Throwable throwable) {
        throwable.printStackTrace();
    }

    @Override
    public void onComplete() {
        System.out.println("Done");
    }
}

// main 클래스
public class FlowTest {
    public static void main(String args[]) {
        List<String> items = List.of("1", "2", "3", "4", "5", "6", "7", "8", "9", "10");
        SubmissionPublisher<String> publisher = new SubmissionPublisher<>();
        publisher.subscribe(new MySubscriber<>());
        items.forEach(s -> {
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            publisher.submit(s);
        });
        publisher.close();
    }
}

실행 결과

1
2
3
4
5
6
7
8
9
10
Done

코드 동작 원리

이 예제의 핵심은 백프레셔(backpressure) 처리 방식입니다. 구독자는 onSubscribe() 시점에 request(1)을 호출해 처음에 단 하나의 항목만 요청하고, onNext()에서 항목을 받을 때마다 다시 request(1)을 호출합니다. 이렇게 하면 구독자가 처리할 준비가 된 만큼만 데이터를 받게 되어, 과부하 없이 안정적인 비동기 데이터 처리가 가능합니다.

모든 항목의 발행이 끝나면 publisher.close()가 호출되고, 이어서 구독자의 onComplete() 콜백이 실행되어 마지막에 "Done"이 출력됩니다.