Java
 Computer >> コンピューター >  >> プログラミング >> Java

Java 9のSubmissionPublisherクラスとは?実装方法とReactive Streamsの基礎を解説

Java 9におけるReactive Streamsの概要

Java 9からは、PublisherSubscriberSubscriptionProcessorという4つのコアインターフェースと、Publisherインターフェースを実装した具象クラスSubmissionPublisherが導入され、Reactive Streams(リアクティブストリーム)を標準APIとして扱えるようになりました。

各インターフェースはそれぞれ異なる役割を持ち、Reactive Streamsの設計原則に対応しています。

  • Publisher:データ(アイテム)を生成し、サブスクライバーへ配信する側の役割を担います。
  • Subscriber:配信されたアイテムを受け取り、処理する側の役割を担います。
  • Subscription:PublisherとSubscriberの関係を管理し、request()によるバックプレッシャー(受信量の制御)を実現します。
  • Processor:PublisherとSubscriberの両方の性質を持ち、データの中継や変換を行います。

SubmissionPublisherクラスのsubmit()メソッドを呼び出すことで、指定したアイテムを現在登録されているすべてのサブスクライバーに非同期で配信できます。

構文(Syntax)

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

SubmissionPublisherはジェネリッククラスであり、AutoCloseableも実装しているため、try-with-resources文との相性も良い点が特徴です。

実装例

以下のコードは、SubmissionPublisherクラスを使った具体的な実装例です。カスタムサブスクライバー「MySubscriber」を定義し、複数のサブスクライバーをパブリッシャーに登録してデータを配信します。

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();
    }
}

コードの解説

  1. MySubscriberクラス:Subscriber<T>インターフェースを実装し、onSubscribe()・onNext()・onError()・onComplete()の4つのコールバックメソッドをオーバーライドしています。
  2. onSubscribe():購読開始時に呼ばれ、Subscriptionインスタンスを保持した後、subscription.request(1)で最初の1件を要求します。
  3. onNext():アイテムを受信するたびに呼ばれ、再度request(1)を呼び出すことで次のアイテムを要求します。これがバックプレッシャーの基本的な動作です。
  4. submit():パブリッシャーがアイテムを発行すると、その時点で登録済みの全サブスクライバーに非同期に配信されます。
  5. close():パブリッシャーを閉じると、全サブスクライバーのonComplete()が呼び出され、ストリームが正常に完了します。

実行結果(Output)

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

実行結果からわかるように、配信は非同期に行われるため、出力順序は実行のたびに変わる可能性があります。また、最後に登録された「Her」サブスクライバーは、それまでに発行されたアイテムを後からまとめて受け取っている点にも注目してください。これは、submit()の呼び出し時点でまだ購読していなかったサブスクライバーには配信されず、購読後に新たに発行された分から処理が始まるためです。

  1. JavaでJPanelのpaintComponent()メソッドを実装する方法をわかりやすく解説

    JPanelはJavaにおける軽量コンテナであり、それ自体は不可視(インビジブル)なコンポーネントとして扱われます。デフォルトのレイアウトにはFlowLayoutが採用されています。JPanelを生成した後は、Containerクラスから継承されたadd()メソッドを呼び出すことで、ボタンやラベルなどの他のコンポーネントをJPanelオブジェクトに自由に追加できます。 paintComponent()メソッドとは このメソッドは、単なる背景色の描画にとどまらず、JPanel上に独自の図形や画像などを描画したい場合に必要となります。paintComponent()メソッドはすでにJPanelクラ

  2. JavaのJToggleButton実装ガイド:ON/OFF切替ボタンの作り方を解説

    JToggleButtonとは JToggleButtonはAbstractButtonを拡張したクラスで、クリックするたびにONとOFFが切り替わるトグルボタンを実現するために使用されます。通常のボタンと異なり、押した状態を保持できるのが特徴です。 JToggleButtonの主な特徴 最初に押されたときは押し込まれた状態のままとなり、もう一度押してはじめて元の状態(押されていない状態)に戻ります。 ボタンが押されるたびにActionEventが発生します。 さらに、JToggleButtonはItemEventも発生させることができます。このイベントは、選択状態という概念を持つコンポー