MQTT サポート
Spring Integration は、メッセージキューテレメトリートランスポート(MQTT)プロトコルをサポートする受信および送信チャネルアダプターを提供します。
この依存関係をプロジェクトに含める必要があります。
現在の実装では、Eclipse Paho MQTT クライアント (英語) ライブラリを使用しています。
両方のアダプターの構成は、DefaultMqttPahoClientFactory を使用して実現されます。構成オプションの詳細については、Paho のドキュメントを参照してください。
ファクトリ自体に(非推奨)オプションを設定する代わりに、MqttConnectOptions オブジェクトを構成してファクトリに注入することをお勧めします。 |
受信(メッセージ駆動型)チャネルアダプター
受信チャネルアダプターは MqttPahoMessageDrivenChannelAdapter によって実装されます。便宜上、名前空間を使用して構成できます。最小構成は次のとおりです。
<bean id="clientFactory"
class="org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory">
<property name="connectionOptions">
<bean class="org.eclipse.paho.client.mqttv3.MqttConnectOptions">
<property name="userName" value="${mqtt.username}"/>
<property name="password" value="${mqtt.password}"/>
</bean>
</property>
</bean>
<int-mqtt:message-driven-channel-adapter id="mqttInbound"
client-id="${mqtt.default.client.id}.src"
url="${mqtt.url}"
topics="sometopic"
client-factory="clientFactory"
channel="output"/>次のリストは、使用可能な属性を示しています。
<int-mqtt:message-driven-channel-adapter id="oneTopicAdapter"
client-id="foo" (1)
url="tcp://localhost:1883" (2)
topics="bar,baz" (3)
qos="1,2" (4)
converter="myConverter" (5)
client-factory="clientFactory" (6)
send-timeout="123" (7)
error-channel="errors" (8)
recovery-interval="10000" (9)
manual-acks="false" (10)
channel="out" />| 1 | クライアント ID。 |
| 2 | ブローカー URL。 |
| 3 | このアダプターがメッセージを受信するトピックのコンマ区切りリスト。 |
| 4 | QoS 値のコンマ区切りリスト。すべてのトピックに適用される単一の値または各トピックの値にすることができます(この場合、リストは同じ長さでなければなりません)。 |
| 5 | MqttMessageConverter (オプション)。デフォルトでは、デフォルトの DefaultPahoMessageConverter は、次のヘッダーを持つ String ペイロードを持つメッセージを生成します。
|
| 6 | クライアントファクトリ。 |
| 7 | 送信タイムアウト。これは、チャネルがブロックする可能性がある場合にのみ適用されます(現在いっぱいの境界付き QueueChannel など)。 |
| 8 | エラーチャネル。ダウンストリーム例外は、ErrorMessage でこのチャネルに送信されます(指定されている場合)。ペイロードは、失敗したメッセージと原因を含む MessagingException です。 |
| 9 | 回復間隔。これは、アダプターが障害後に再接続を試行する間隔を制御します。デフォルトは 10000ms (10 秒)です。 |
| 10 | 確認応答モード。手動で確認する場合は true に設定します。 |
バージョン 4.1 以降、URL を省略できます。代わりに、DefaultMqttPahoClientFactory の serverURIs プロパティでサーバー URI を提供できます。これにより、たとえば、高可用性(HA)クラスターへの接続が可能になります。 |
バージョン 4.2.2 以降、アダプターがトピックを正常にサブスクライブすると、MqttSubscribedEvent が公開されます。接続またはサブスクリプションが失敗すると、MqttConnectionFailedEvent イベントが発行されます。ApplicationListener を実装する Bean がこれらのイベントを受信できます。
また、recoveryInterval という新しいプロパティは、アダプターが障害後に再接続を試みる間隔を制御します。デフォルトは 10000ms (10 秒)です。
バージョン 4.2.3 より前は、クライアントはアダプターが停止したときに常にサブスクライブ解除されました。クライアントの QOS が 0 より大きい場合、アダプターが停止している間に到着したメッセージが次回の開始時に配信されるように、サブスクリプションをアクティブに保つ必要があるため、これは間違っていました。これには、クライアントファクトリの バージョン 4.2.3 以降、 この動作は、ファクトリで 4.2.3 より前の動作に戻すには、 |
バージョン 5.0 以降、 |
実行時のトピックの追加と削除
バージョン 4.1 以降、アダプターがサブスクライブするトピックをプログラムで変更できます。Spring Integration は、addTopic() および removeTopic() メソッドを提供します。トピックを追加するときに、オプションで QoS を指定できます(デフォルト: 1)。適切なメッセージを適切なペイロードで <control-bus/> に送信して、トピックを変更することもできます(例: "myMqttAdapter.addTopic('foo', 1)")。
アダプターを停止して開始しても、トピックリストには影響しません(構成の元の設定に戻りません)。変更は、アプリケーションコンテキストのライフサイクルを超えて保持されません。新しいアプリケーションコンテキストは、構成された設定に戻ります。
アダプターが停止している(またはブローカーから切断されている)間にトピックを変更すると、次に接続が確立されたときに有効になります。
手動 ACK
バージョン 5.3 以降、manualAcks プロパティを true に設定できます。多くの場合、配信を非同期的に確認するために使用されます。true に設定すると、ヘッダー(IntegrationMessageHeaderAccessor.ACKNOWLEDGMENT_CALLBACK)がメッセージに追加され、値は SimpleAcknowledgment になります。配信を完了するには、acknowledge() メソッドを呼び出す必要があります。詳細については、IMqppClient setManualAcks()、messageArrivedComplete() の Javadoc を参照してください。便宜上、ヘッダーアクセサーが提供されています。
StaticMessageHeaderAccessor.acknowledgment(someMessage).acknowledge();Java 構成を使用した構成
次の Spring Boot アプリケーションは、Java 構成で受信アダプターを構成する方法の例を示しています。
@SpringBootApplication
public class MqttJavaApplication {
public static void main(String[] args) {
new SpringApplicationBuilder(MqttJavaApplication.class)
.web(false)
.run(args);
}
@Bean
public MessageChannel mqttInputChannel() {
return new DirectChannel();
}
@Bean
public MessageProducer inbound() {
MqttPahoMessageDrivenChannelAdapter adapter =
new MqttPahoMessageDrivenChannelAdapter("tcp://localhost:1883", "testClient",
"topic1", "topic2");
adapter.setCompletionTimeout(5000);
adapter.setConverter(new DefaultPahoMessageConverter());
adapter.setQos(1);
adapter.setOutputChannel(mqttInputChannel());
return adapter;
}
@Bean
@ServiceActivator(inputChannel = "mqttInputChannel")
public MessageHandler handler() {
return new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
System.out.println(message.getPayload());
}
};
}
}Java DSL を使用した構成
次の Spring Boot アプリケーションは、Java DSL を使用して受信アダプターを構成する例を示しています。
@SpringBootApplication
public class MqttJavaApplication {
public static void main(String[] args) {
new SpringApplicationBuilder(MqttJavaApplication.class)
.web(false)
.run(args);
}
@Bean
public IntegrationFlow mqttInbound() {
return IntegrationFlows.from(
new MqttPahoMessageDrivenChannelAdapter("tcp://localhost:1883",
"testClient", "topic1", "topic2");)
.handle(m -> System.out.println(m.getPayload()))
.get();
}
}送信チャネルアダプター
発信チャネルアダプターは MqttPahoMessageHandler によって実装され、ConsumerEndpoint にラップされています。便宜上、名前空間を使用して構成できます。
バージョン 4.1 から、アダプターは非同期送信操作をサポートし、配信が確認されるまでブロッキングを回避します。必要に応じて、アプリケーションが配信を確認できるように、アプリケーションイベントを発行できます。
次のリストは、送信チャネルアダプターで使用可能な属性を示しています。
<int-mqtt:outbound-channel-adapter id="withConverter"
client-id="foo" (1)
url="tcp://localhost:1883" (2)
converter="myConverter" (3)
client-factory="clientFactory" (4)
default-qos="1" (5)
qos-expression="" (6)
default-retained="true" (7)
retained-expression="" (8)
default-topic="bar" (9)
topic-expression="" (10)
async="false" (11)
async-events="false" (12)
channel="target" />| 1 | クライアント ID。 |
| 2 | ブローカー URL。 |
| 3 | MqttMessageConverter (オプション)。デフォルトの DefaultPahoMessageConverter は、次のヘッダーを認識します。
|
| 4 | クライアントファクトリ。 |
| 5 | デフォルトのサービス品質。mqtt_qos ヘッダーが見つからない場合、または qos-expression が null を返す場合に使用されます。カスタム converter を提供する場合は使用されません。 |
| 6 | QoS を決定するために評価する式。デフォルトは headers[mqtt_qos] です。 |
| 7 | 保持フラグのデフォルト値。mqtt_retained ヘッダーが見つからない場合に使用されます。カスタム converter が提供されている場合は使用されません。 |
| 8 | 保持されたブール値を決定するために評価する式。デフォルトは headers[mqtt_retained] です。 |
| 9 | メッセージの送信先のデフォルトトピック(mqtt_topic ヘッダーが見つからない場合に使用)。 |
| 10 | 宛先トピックを決定するために評価する式。デフォルトは headers['mqtt_topic'] です。 |
| 11 | true の場合、呼び出し元はブロックしません。むしろ、メッセージが送信されると配信確認を待ちます。デフォルトは false です(配信が確認されるまで送信はブロックされます)。 |
| 12 | async と async-events が両方とも true である場合、MqttMessageSentEvent が発行されます(イベントを参照)。メッセージ、トピック、クライアントライブラリによって生成された messageId、clientId、clientInstance (クライアントが接続されるたびにインクリメントされます)が含まれます。配信がクライアントライブラリによって確認されると、MqttMessageDeliveredEvent が発行されます。messageId、clientId、clientInstance が含まれているため、配信と送信を関連付けることができます。ApplicationListener またはイベント受信チャネルアダプターは、これらのイベントを受信できます。MqttMessageDeliveredEvent が MqttMessageSentEvent の前に受信される可能性があることに注意してください。デフォルトは false です。 |
バージョン 4.1 以降、URL は省略できます。代わりに、DefaultMqttPahoClientFactory の serverURIs プロパティでサーバー URI を提供できます。これにより、たとえば、高可用性(HA)クラスターへの接続が可能になります。 |
Java 構成を使用した構成
次の Spring Boot アプリケーションは、Java 構成で送信アダプターを構成する方法の例を示しています。
@SpringBootApplication
@IntegrationComponentScan
public class MqttJavaApplication {
public static void main(String[] args) {
ConfigurableApplicationContext context =
new SpringApplicationBuilder(MqttJavaApplication.class)
.web(false)
.run(args);
MyGateway gateway = context.getBean(MyGateway.class);
gateway.sendToMqtt("foo");
}
@Bean
public MqttPahoClientFactory mqttClientFactory() {
DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
MqttConnectOptions options = new MqttConnectOptions();
options.setServerURIs(new String[] { "tcp://host1:1883", "tcp://host2:1883" });
options.setUserName("username");
options.setPassword("password".toCharArray());
factory.setConnectionOptions(options);
return factory;
}
@Bean
@ServiceActivator(inputChannel = "mqttOutboundChannel")
public MessageHandler mqttOutbound() {
MqttPahoMessageHandler messageHandler =
new MqttPahoMessageHandler("testClient", mqttClientFactory());
messageHandler.setAsync(true);
messageHandler.setDefaultTopic("testTopic");
return messageHandler;
}
@Bean
public MessageChannel mqttOutboundChannel() {
return new DirectChannel();
}
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MyGateway {
void sendToMqtt(String data);
}
}Java DSL を使用した構成
次の Spring Boot アプリケーションは、Java DSL を使用して送信アダプターを構成する例を示しています。
@SpringBootApplication
public class MqttJavaApplication {
public static void main(String[] args) {
new SpringApplicationBuilder(MqttJavaApplication.class)
.web(false)
.run(args);
}
@Bean
public IntegrationFlow mqttOutboundFlow() {
return f -> f.handle(new MqttPahoMessageHandler("tcp://host1:1883", "someMqttClient"));
}
}イベント
特定のアプリケーションイベントは、アダプターによって発行されます。
MqttConnectionFailedEvent- 接続に失敗した場合、またはその後接続が失われた場合に、両方のアダプターによって公開されます。MqttMessageSentEvent- 非同期モードで実行されている場合、メッセージが送信されたときに送信アダプターによって公開されます。MqttMessageDeliveredEvent- 非同期モードで実行している場合に、クライアントがメッセージが配信されたことを示したときに送信アダプターによって公開されます。MqttSubscribedEvent- トピックをサブスクライブした後、受信アダプターによって公開されます。
これらのイベントは、ApplicationListener<MqttIntegrationEvent> または @EventListener メソッドで受信できます。