微信公众号搜"智元新知"关注
微信扫一扫可直接关注哦!

Java:反应式MQTT-并行订阅

如何解决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 举报,一经查实,本站将立刻删除。