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

Java 9でSubscriberインターフェースを実装する方法を解説

Java 9では、Reactive Streams(リアクティブストリーム)を標準ライブラリで扱えるようにするため、いくつかの新しいインターフェースが導入されました。主な構成要素はPublisherSubscriberSubscriptionの各インターフェースと、Publisherインターフェースを実装したSubmissionPublisherクラスです。それぞれがReactive Streamsの原則に応じて異なる役割を担います。

Subscriberインターフェースを使うことで、パブリッシャー(発行者)が配信するデータを購読できます。実際に利用するには、Subscriberインターフェースを実装するクラスを作成し、抽象メソッドの処理内容を定義する必要があります。

Flow.Subscriberインターフェースの主なメソッド

  • onComplete():Publisherオブジェクトがその役割を完了した際に呼び出されます。
  • onError():Publisher側でエラーが発生し、それがSubscriberに通知された際に呼び出されます。
  • onNext():Publisherに新しいデータが発生し、すべてのSubscriberへ通知されるたびに呼び出されます。
  • onSubscribe():PublisherにSubscriberが登録された際に呼び出されます。

実装例

以下のコードは、1から10までの整数を発行するSubmissionPublisherに対して、独自のSubscriberを実装して購読する例です。

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

コードのポイント

  • バックプレッシャー(背圧)の制御:onSubscribe()およびonNext()内でsubscription.request(1)を呼び出すことで、一度に1件ずつデータを要求しています。これにより、処理能力を超えたデータの流入を防ぐことができます。
  • 完了の検知:onComplete()が呼ばれたタイミングでisDoneフラグを立て、mainスレッドはそのフラグを監視しながら処理の終了を待機します。
  • 非同期処理:SubmissionPublisherはデフォルトで別スレッド上でアイテムを配信するため、publish処理と購読処理が非同期に動作します。

このように、Java 9のFlow APIを利用することで、外部ライブラリに依存せずにReactive Streamsの基本的な仕組み(Publisher/Subscriberモデルとバックプレッシャー)を簡単に試すことができます。

  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も発生させることができます。このイベントは、選択状態という概念を持つコンポー