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 (将#修改为@)