如何解决Java:反应式MQTT-并行订阅
我想与hivemq-mqtt-client订阅一个mqtt主题。它支持反应式api,但我并不十分熟悉。
我想为每个输入消息使用一个自己的轻量级线程。 目前,当输入法被阻止时,我无法收到更多消息:
Mqtt5ClientBuilder builder = this.getClientBuilder(host);
try {
setupSslConfig(builder,caPath);
} catch (Exception e) {
e.printstacktrace();
}
this.mqttClient = builder.buildrx();
this.mqttClient.connect();
Single<Mqtt5ConnAck> connAckSingle = this.mqttClient.connect();
// Filter for the init/ topic
FlowableWithSingle<Mqtt5Publish,Mqtt5SubAck> subAckAndMatchingPublishes = this.mqttClient.subscribePublishesWith()
.topicFilter("init/").qos(MqttQos.AT_LEAST_ONCE)
.applySubscribe();
// Connection handler (debug outputs)
Completable connectScenario = connAckSingle
.doOnSuccess(connAck -> System.out.println("Connected," + connAck.getReasonCode()))
.doOnError(throwable -> System.out.println("Connection Failed," + throwable.getMessage()))
.ignoreElement();
// Connection handler (debug outputs)
Completable subscribeScenario = subAckAndMatchingPublishes
.doOnSingle(subAck -> System.out.println("Subscribed," + subAck.getReasonCodes()))
.doOnNext(publish -> this.processMessage(publish.getTopic().toString(),publish.getPayloadAsBytes()))
.ignoreElements();
connectScenario.andThen(subscribeScenario).blockingAwait();
public void processMessage(String topic,byte[] raw) {
System.out.println("CHECK");
while(true) {
}
}
如果我收到有关特定主题的消息,则会在控制台中阅读一次“ CHECK”。然后,侦听mqtt代理的线程将被阻止...
我知道我可以在processMessage
内生成一个织机或线程,但是我希望有一种使用反应式编程的方法。
T
版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。