You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus中GraphQL订阅接收所有MQTT主题消息的问题求助

问题诊断与解决方案

你的问题核心在于没有按MQTT主题过滤消息,且重复创建处理器与MQTT订阅绑定,导致所有GraphQL订阅都能收到MQTT客户端收到的全部消息。具体来说:

  1. 每次调用subscribe(String topic)都会新建一个BroadcastProcessor,同时给MQTT客户端设置一个全局的publishHandler——这个handler没有检查消息所属的主题,只要MQTT客户端收到消息(不管是哪个主题的),就会发送给当前创建的processor。
  2. 多次调用client.subscribe(topic, 2)会让MQTT客户端订阅多个主题,导致客户端能收到所有这些主题的消息,而每个processor都会收到所有这些消息(因为handler没有主题过滤逻辑)。

修复方案:主题映射+全局消息分发

我们需要维护一个主题到BroadcastProcessor的映射表,让每个MQTT主题对应唯一的处理器,同时只设置一次全局的publishHandler来按主题分发消息。这样既能避免重复订阅MQTT主题,又能保证只有订阅对应主题的GraphQL客户端收到消息。

步骤1:重构MQTT订阅逻辑

首先创建一个单例的管理器类(或者在你的MQTT客户端类中维护状态),负责维护主题与处理器的映射,以及全局的消息分发:

import io.smallrye.mutiny.Multi;
import io.smallrye.mutiny.subscription.BroadcastProcessor;
import io.vertx.mqtt.MqttClient;
import io.vertx.mqtt.messages.MqttPublishMessage;

import javax.enterprise.context.ApplicationScoped;
import javax.inject.Inject;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

@ApplicationScoped
public class MqttSubscriptionManager {

    private final MqttClient client;
    private final Map<String, BroadcastProcessor<String>> topicProcessors = new ConcurrentHashMap<>();

    @Inject
    public MqttSubscriptionManager(MqttClient client) {
        this.client = client;
        // 设置一次全局的publishHandler,负责按主题分发消息
        this.client.publishHandler(this::dispatchMessageToProcessor);
    }

    private void dispatchMessageToProcessor(MqttPublishMessage message) {
        String topic = message.topicName();
        String payload = message.payload().toString();
        System.out.println("Received message on topic: " + topic + ", payload: " + payload);
        
        // 只把消息发送给对应主题的processor
        BroadcastProcessor<String> processor = topicProcessors.get(topic);
        if (processor != null && !processor.isTerminated()) {
            processor.onNext(payload);
        }
    }

    public Multi<String> subscribe(String topic) {
        return topicProcessors.computeIfAbsent(topic, this::createProcessorForTopic);
    }

    private BroadcastProcessor<String> createProcessorForTopic(String topic) {
        BroadcastProcessor<String> processor = BroadcastProcessor.create();
        // 订阅MQTT主题(只订阅一次)
        client.subscribe(topic, 2)
                .onFailure().invoke(err -> {
                    System.err.println("Failed to subscribe to topic: " + topic + ", error: " + err.getMessage());
                    processor.onError(err);
                });
        
        // 当processor没有订阅者时,清理资源并取消MQTT订阅
        processor.onTermination().invoke(() -> {
            topicProcessors.remove(topic);
            client.unsubscribe(topic)
                    .onSuccess(v -> System.out.println("Unsubscribed from topic: " + topic))
                    .onFailure().invoke(err -> System.err.println("Failed to unsubscribe from topic: " + topic + ", error: " + err.getMessage()));
        });
        
        return processor;
    }
}

步骤2:更新GraphQL订阅实现

现在你的GraphQL订阅方法只需要调用管理器的subscribe方法即可:

@Subscription
public Multi<String> messageCreated(String topic) {
    return mqttSubscriptionManager.subscribe(topic).toHotStream();
}

关键改进点说明

  • 主题映射表:用ConcurrentHashMap维护主题与processor的对应关系,确保同一个主题的多个GraphQL订阅共享同一个processor,避免重复订阅MQTT主题。
  • 全局消息分发:只设置一次publishHandler,根据消息的主题找到对应的processor,只转发该主题的消息,解决了所有订阅收到全部消息的问题。
  • 资源清理:当processor没有订阅者时(比如GraphQL客户端断开连接),自动取消对应的MQTT主题订阅,并从映射表中移除processor,避免资源泄漏。
  • 错误处理:添加了订阅失败的错误处理,确保GraphQL订阅能收到错误通知。

这样修改后,每个GraphQL订阅只会收到自己订阅的MQTT主题的消息,同时避免了重复订阅与资源浪费。

内容的提问来源于stack exchange,提问作者dna

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 23:37:29