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

Java 9에서 Subscriber(구독자) 인터페이스 구현하는 방법

Java 9Reactive 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

코드 설명

위 코드의 동작 흐름은 다음과 같습니다.

  1. 구독 등록: publisher.subscribe(subscriber)가 호출되면 Subscriber의 onSubscribe() 메서드가 실행되고, 이때 subscription.request(1)을 통해 처음 한 개의 데이터를 요청합니다.
  2. 데이터 처리: Publisher가 submit()으로 데이터를 발행하면 onNext()가 호출되어 각 항목을 처리한 뒤, 다시 request(1)로 다음 데이터를 요청하는 백프레셔(backpressure) 방식으로 동작합니다.
  3. 완료 처리: publisher.close()가 호출되면 스트림이 종료되고, onComplete() 메서드가 실행되어 모든 처리가 끝났음을 알립니다.

이처럼 Java 9의 Flow API를 활용하면 비동기 데이터 스트림을 안정적으로 처리할 수 있으며, 요청 기반(request-based)의 백프레셔 제어를 통해 생산자와 소비자 간의 데이터 흐름을 효율적으로 관리할 수 있습니다.