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

Java 9에서 Flow.Publisher 인터페이스 구현하는 방법

Publisher 인터페이스는 무제한 개수의 순차적 요소를 제공하는 공급자 역할을 하며, Subscriber(구독자)로부터 전달받은 수요(demand)에 따라 요소들을 발행(publish)합니다. Java 9부터 표준 라이브러리에 포함된 java.util.concurrent.Flow API는 리액티브 스트림(Reactive Streams) 사양을 기반으로 하며, 비동기 환경에서 안정적인 데이터 흐름 처리를 지원합니다.

Publisher.subscribe(Subscriber) 메서드가 호출되면 Subscriber의 메서드들은 일반적으로 다음과 같은 순서로 호출됩니다.

  • 먼저 onSubscribe() 메서드가 한 번 호출됩니다.
  • 이후 Subscriber가 요청한 만큼 onNext() 메서드가 반복적으로 호출됩니다.
  • 처리 중 오류가 발생하면 onError()가, 더 이상 발행할 요소가 없으면 onComplete()가 호출되며, 이는 Subscription이 취소되지 않은 경우에만 유효합니다.

문법(Syntax)

public interface Publisher<T> {
    public void subscribe(Subscriber<? super T> s);
}

Publisher 인터페이스에는 단 하나의 추상 메서드인 subscribe()만 존재하기 때문에 함수형 인터페이스로 간주할 수 있으며, 구현체는 이 메서드를 통해 Subscriber와 연결됩니다.

예제 코드

다음은 1부터 지정한 개수까지의 정수를 발행하는 간단한 SimplePublisher 구현 예제입니다.

import java.util.concurrent.*;
import java.util.*;
import java.util.stream.*;

class SimplePublisher implements Flow.Publisher<Integer> {
    private final Iterator<Integer> iterator;

    SimplePublisher(int count) {
        this.iterator = IntStream.rangeClosed(1, count).iterator();
    }

    @Override
    public void subscribe(Flow.Subscriber<? super Integer> subscriber) {
        iterator.forEachRemaining(subscriber::onNext);
        subscriber.onComplete();
    }
}

public class SimplePublisherImplTest {
    public static void main(String args[]) {
        new SimplePublisher(10).subscribe(new Flow.Subscriber<>() {
            @Override
            public void onSubscribe(Flow.Subscription subscription) {
            }
            @Override
            public void onNext(Integer item) {
                System.out.println("item = [" + item + "]");
            }
            @Override
            public void onError(Throwable throwable) {
            }
            @Override
            public void onComplete() {
                System.out.println("complete");
            }
        });
    }
}

코드 설명

  • SimplePublisher 클래스: IntStream.rangeClosed(1, count)로 생성한 이터레이터를 내부에 보관하여, subscribe() 호출 시 요소를 순차적으로 발행합니다.
  • subscribe() 메서드: forEachRemaining()으로 남은 모든 요소를 subscriber::onNext에 전달한 뒤, 마지막에 onComplete()를 호출해 발행 종료를 알립니다.
  • 익명 Subscriber: Flow.Subscriber의 네 가지 콜백 메서드(onSubscribe, onNext, onError, onComplete)를 구현하여 발행되는 각 항목을 콘솔에 출력합니다.

실행 결과

item = [1]
item = [2]
item = [3]
item = [4]
item = [5]
item = [6]
item = [7]
item = [8]
item = [9]
item = [10]
complete

참고 사항

위 예제는 개념 이해를 돕기 위한 최소 구현으로, Subscriber의 request(n) 호출을 통한 수요(demand) 관리, 즉 백프레셔(backpressure) 처리는 반영되어 있지 않습니다. 실제 운영 환경에서는 Subscription 객체를 활용해 Subscriber가 요청한 만큼만 데이터를 발행하도록 구현하는 것이 권장됩니다. 이러한 규격을 준수하는 검증 도구로는 SubmissionPublisher나 리액티브 스트림 TCK(Technology Compatibility Kit)를 활용할 수 있습니다.