Quarkus中GraphQL订阅接收所有MQTT主题消息的问题求助
问题诊断与解决方案
你的问题核心在于没有按MQTT主题过滤消息,且重复创建处理器与MQTT订阅绑定,导致所有GraphQL订阅都能收到MQTT客户端收到的全部消息。具体来说:
- 每次调用
subscribe(String topic)都会新建一个BroadcastProcessor,同时给MQTT客户端设置一个全局的publishHandler——这个handler没有检查消息所属的主题,只要MQTT客户端收到消息(不管是哪个主题的),就会发送给当前创建的processor。 - 多次调用
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
相关产品推荐
相关产品推荐

