このバージョンはまだ開発中であり、まだ安定しているとは考えられていません。最新の安定バージョンについては、spring-cloud-stream 5.0.3 を使用してください。

結合の視覚化と制御

Spring Cloud Stream は、アクチュエーターエンドポイントを介したバインディングの視覚化と制御、およびプログラムによる方法をサポートしています。

プログラム的な方法

バージョン 3.1 以降、Bean として登録されている org.springframework.cloud.stream.binding.BindingsLifecycleController を公開しており、注入されると、個々のバインディングのライフサイクルを制御するために使用できます。

ex: テストケースの 1 つのフラグメントを調べます。ご覧のとおり、Spring アプリケーションコンテキストから BindingsLifecycleController を取得し、個々のメソッドを実行して echo-in-0 バインディングのライフサイクルを制御します。

BindingsLifecycleController bindingsController = context.getBean(BindingsLifecycleController.class);
Binding binding = bindingsController.queryState("echo-in-0");
assertThat(binding.isRunning()).isTrue();
bindingsController.changeState("echo-in-0", State.STOPPED);
//Alternative way of changing state. For convenience we expose start/stop and pause/resume operations.
//bindingsController.stop("echo-in-0")
assertThat(binding.isRunning()).isFalse();

新しいバインディングを定義し、既存のバインディングを管理する

さらに、BindingsLifecycleController を使用してバージョン 4.2 を開始すると、コンシューマーおよびプロデューサーの構成プロパティにアクセスして値をより動的に管理することで、新しいバインディングを定義したり、既存のバインディング構成を変更したりできるようになります。

次に例を示します。

新しい入力バインディングを定義するには、BindingsLifecycleController.createInputBinding(..) メソッド(下記参照)を呼び出します。createOutputBinding(..) メソッドに相当するメソッドもあります。

DefaultBinderFactory binderFactory = context.getBean(DefaultBinderFactory.class);
Object binder = binderFactory.getBinder("rabbit", MessageChannel.class);

String inputName = "test-output-binding";

BindingsLifecycleController controller = context.getBean(BindingsLifecycleController.class);

BindingProperties bindingProperties = new BindingProperties();
bindingProperties.setDestination("myDest");
bindingProperties.setGroup("myGroup"); // must define group for input binding
ConsumerProperties consumerProperties = new ConsumerProperties();
bindingProperties.setConsumer(consumerProperties);

RabbitConsumerProperties extendedConsumerProperties = controller.createInputBinding(inputName, "rabbit", bindingProperties);

その後、getExtensionProperties(..) メソッドを呼び出してそのプロパティを管理できます。

KafkaConsumerProperties properties = controller.getExtensionProperties("test-input-binding”);
関数定義から派生したバインディング名とは異なり、明示的に定義されたバインディングは実際の関数によってサポートされていないため、in-0/out-0 サフィックスは付きません。

getExtensionProperties(..) 操作は、構成プロパティクラスの適切な型を取得するために定義されています。そのため、使用するバインダーとバインディングに応じて、拡張プロパティを適切な型に安全にキャストできます。この例では、KafkaConsumerProperties プロパティです。

変更するプロパティの種類によっては、変更を有効にするためにバインディングを再起動する必要がある場合があります。(先ほど見たお尻)

アクチュエーター

アクチュエーターと Web はオプションであるため、最初に Web 依存関係の 1 つを追加し、アクチュエーター依存関係を手動で追加する必要があります。次の例は、Web フレームワークの依存関係を追加する方法を示しています。

<dependency>
     <groupId>org.springframework.boot</groupId>
     <artifactId>spring-boot-starter-web</artifactId>
</dependency>

次の例は、WebFlux フレームワークの依存関係を追加する方法を示しています。

<dependency>
       <groupId>org.springframework.boot</groupId>
       <artifactId>spring-boot-starter-webflux</artifactId>
</dependency>

次のように、アクチュエーターの依存関係を追加できます。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
Cloud Foundry で Spring Cloud Stream 2.0 アプリを実行するには、クラスパスに spring-boot-starter-web と spring-boot-starter-actuator を追加する必要があります。そうしないと、ヘルスチェックの失敗によりアプリケーションが起動しません。

また、次のプロパティを設定して、bindings アクチュエーターのエンドポイントを有効にする必要があります: --management.endpoints.web.exposure.include=bindings。

これらの前提条件が満たされたら。アプリケーションの起動時に、ログに次の情報が表示されます。

: Mapped "{[/actuator/bindings/{name}],methods=[POST]. . .
: Mapped "{[/actuator/bindings],methods=[GET]. . .
: Mapped "{[/actuator/bindings/{name}],methods=[GET]. . .

現在のバインディングを視覚化するには、次の URL にアクセスします: <host>:<port>/actuator/bindings (英語)

別の方法として、単一のバインディングを表示するには、次のような URL のいずれかにアクセスします。<host>:<port>/actuator/bindings/<bindingName> (英語) ;

次の例に示すように、JSON として state 引数を指定しながら、同じ URL に投稿することで、個々のバインディングを停止、開始、一時停止、再開することもできます。

curl -d '{"state":"STOPPED"}' -H "Content-Type: application/json" -X POST http://<host>:<port>/actuator/bindings/myBindingName
curl -d '{"state":"STARTED"}' -H "Content-Type: application/json" -X POST http://<host>:<port>/actuator/bindings/myBindingName
curl -d '{"state":"PAUSED"}' -H "Content-Type: application/json" -X POST http://<host>:<port>/actuator/bindings/myBindingName
curl -d '{"state":"RESUMED"}' -H "Content-Type: application/json" -X POST http://<host>:<port>/actuator/bindings/myBindingName
PAUSED および RESUMED は、対応するバインダーとその基礎となるテクノロジーがサポートしている場合にのみ機能します。それ以外の場合は、ログに警告メッセージが表示されます。現在、Kafka および [Solace]( github.com/SolaceProducts/solace-spring-cloud/tree/master/solace-spring-cloud-starters/solace-spring-cloud-stream-starter#consumer-bindings-pauseresume (英語) ) バインダーのみが PAUSED および RESUMED 状態をサポートしています。

機密データのサニタイズ

バインディングアクチュエーターエンドポイントを使用する場合、ユーザー資格情報、SSL キーに関する情報などの機密データをサニタイズすることが重要な場合があります。これを実現するために、エンドユーザーアプリケーションは、アプリケーションで Bean として Spring Boot の SanitizingFunction を提供できます。Apache Kafka の sasl.jaas.config プロパティに値を提供するときにデータをスクランブルする例を次に示します。

@Bean
public SanitizingFunction sanitizingFunction() {
	return sanitizableData -> {
		if (sanitizableData.getKey().equals("sasl.jaas.config")) {
			return sanitizableData.withValue("data-scrambled!!");
		}
		else {
			return sanitizableData;
		}
	};
}