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

Java 9 SubmissionPublisher 클래스로 리액티브 스트림 구현하기

Java 9 리액티브 스트림과 SubmissionPublisher 소개

Java 9부터는 Publisher, Subscriber, Subscription, Processor라는 네 가지 핵심 인터페이스와, 그중 Publisher 인터페이스를 구현한 구체 클래스인 SubmissionPublisher가 표준 라이브러리에 추가되면서 리액티브 스트림(Reactive Streams) 프로그래밍이 가능해졌습니다. 각 인터페이스는 리액티브 스트림의 기본 원칙에 따라 서로 다른 역할을 담당합니다.

SubmissionPublisher 클래스의 submit() 메서드를 호출하면 전달된 아이템을 현재 등록된 모든 구독자에게 비동기적으로 발행할 수 있습니다.

클래스 선언(문법)

public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable

아래 예제에서 SubmissionPublisher 클래스를 실제로 구현하고 사용하는 방법을 살펴보겠습니다.

예제 코드

import java.util.concurrent.Flow.Subscriber;
import java.util.concurrent.Flow.Subscription;
import java.util.concurrent.SubmissionPublisher;

class MySubscriber<T> implements Subscriber<T> {
    private Subscription subscription;
    private String name;

    public MySubscriber(String name) {
        this.name = name;
    }

    @Override
    public void onComplete() {
        System.out.println(name + ": onComplete");
    }

    @Override
    public void onError(Throwable t) {
        System.out.println(name + ": onError");
        t.printStackTrace();
    }

    @Override
    public void onNext(T msg) {
        System.out.println(name + ": " + msg.toString() + " received in onNext");
        subscription.request(1);
    }

    @Override
    public void onSubscribe(Subscription subscription) {
        System.out.println(name + ": onSubscribe");
        this.subscription = subscription;
        subscription.request(1);
    }
}

// 메인 클래스
public class FlowTest {
    public static void main(String args[]) {
        SubmissionPublisher<String> publisher = new SubmissionPublisher<>();
        MySubscriber<String> subscriber = new MySubscriber<>("Mine");
        MySubscriber<String> subscriberYours = new MySubscriber<>("Yours");
        MySubscriber<String> subscriberHis = new MySubscriber<>("His");
        MySubscriber<String> subscriberHers = new MySubscriber<>("Her");

        publisher.subscribe(subscriber);
        publisher.subscribe(subscriberYours);
        publisher.subscribe(subscriberHis);
        publisher.subscribe(subscriberHers);

        publisher.submit("One");
        publisher.submit("Two");
        publisher.submit("Three");
        publisher.submit("Four");
        publisher.submit("Five");

        try {
            Thread.sleep(1000);
        } catch(InterruptedException e) {
            e.printStackTrace();
        }
        publisher.close();
    }
}

코드 설명

  • MySubscriber: Subscriber 인터페이스를 구현한 커스텀 구독자 클래스입니다.
  • onSubscribe(): 구독이 시작될 때 호출되며, Subscription 객체를 저장한 뒤 request(1)로 아이템을 요청합니다.
  • onNext(): 새로운 아이템이 도착할 때마다 호출되며, 처리 후 다음 아이템을 다시 요청합니다.
  • onError(): 처리 과정에서 오류가 발생하면 호출됩니다.
  • onComplete(): 모든 아이템 처리가 정상적으로 완료되면 호출됩니다.

구독자가 request(1)을 호출하는 방식은 리액티브 스트림의 핵심 개념인 배압(backpressure)을 잘 보여줍니다. 구독자가 자신이 감당할 수 있는 만큼만 아이템을 요청함으로써 데이터 흐름을 능동적으로 제어할 수 있습니다.

실행 결과

Yours: onSubscribe
His: onSubscribe
Mine: onSubscribe
His: One received in onNext
Yours: One received in onNext
Mine: One received in onNext
Yours: Two received in onNext
His: Two received in onNext
Yours: Three received in onNext
Mine: Two received in onNext
Yours: Four received in onNext
His: Three received in onNext
Yours: Five received in onNext
Mine: Three received in onNext
Her: onSubscribe
His: Four received in onNext
Her: One received in onNext
Mine: Four received in onNext
Her: Two received in onNext
His: Five received in onNext
Her: Three received in onNext
Mine: Five received in onNext
Her: Four received in onNext
Her: Five received in onNext
Yours: onComplete
His: onComplete
Mine: onComplete
Her: onComplete

실행 결과에서 볼 수 있듯이, 각 구독자는 발행된 아이템을 개별적으로 수신하며 publisher.close()가 호출되면 모든 구독자에게 onComplete()가 전달되어 스트림이 정상적으로 종료됩니다. 발행과 구독이 비동기적으로 동작하기 때문에 출력 순서는 실행 시점에 따라 달라질 수 있다는 점도 참고하세요.