Java 9でSubscriberインターフェースを実装する方法を解説
Java 9では、Reactive Streams(リアクティブストリーム)を標準ライブラリで扱えるようにするため、いくつかの新しいインターフェースが導入されました。主な構成要素はPublisher、Subscriber、Subscriptionの各インターフェースと、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モデルとバックプレッシャー)を簡単に試すことができます。
-
JavaでJPanelのpaintComponent()メソッドを実装する方法をわかりやすく解説
JPanelはJavaにおける軽量コンテナであり、それ自体は不可視(インビジブル)なコンポーネントとして扱われます。デフォルトのレイアウトにはFlowLayoutが採用されています。JPanelを生成した後は、Containerクラスから継承されたadd()メソッドを呼び出すことで、ボタンやラベルなどの他のコンポーネントをJPanelオブジェクトに自由に追加できます。 paintComponent()メソッドとは このメソッドは、単なる背景色の描画にとどまらず、JPanel上に独自の図形や画像などを描画したい場合に必要となります。paintComponent()メソッドはすでにJPanelクラ
-
JavaのJToggleButton実装ガイド:ON/OFF切替ボタンの作り方を解説
JToggleButtonとは JToggleButtonはAbstractButtonを拡張したクラスで、クリックするたびにONとOFFが切り替わるトグルボタンを実現するために使用されます。通常のボタンと異なり、押した状態を保持できるのが特徴です。 JToggleButtonの主な特徴 最初に押されたときは押し込まれた状態のままとなり、もう一度押してはじめて元の状態(押されていない状態)に戻ります。 ボタンが押されるたびにActionEventが発生します。 さらに、JToggleButtonはItemEventも発生させることができます。このイベントは、選択状態という概念を持つコンポー