Spring Integration MQTTv5共享订阅无法接收消息问题
问题解决:Spring Integration MQTT共享订阅消息无法到达处理器
问题原因
使用共享订阅格式$share/testGroup/foo/test时,MQTT Broker推送的消息主题为原始主题foo/test,而Mqttv5PahoMessageDrivenChannelAdapter默认会以订阅时的完整主题(带$share前缀)作为消息选择器的匹配条件,导致推送的消息主题与选择器不匹配,被过滤掉,无法传递到后续处理器。
解决方案
禁用Mqttv5PahoMessageDrivenChannelAdapter的默认消息选择器,让适配器接收所有Broker推送的对应共享组主题消息。
修改后的代码如下:
@Bean fun clientManager(): ClientManager<IMqttAsyncClient, MqttConnectionOptions> { val connectionOptions = MqttConnectionOptions() connectionOptions.serverURIs = arrayOf("tcp://example.org:1883") val clientManager = Mqttv5ClientManager(connectionOptions, "testClient") clientManager.setPersistence(MqttDefaultFilePersistence()) return clientManager } @Bean fun mqttTestInFlow(clientManager: ClientManager<IMqttAsyncClient, MqttConnectionOptions>): IntegrationFlow { val messageProducer = Mqttv5PahoMessageDrivenChannelAdapter( clientManager, "\$share/testGroup/foo/test", ) // 禁用默认消息选择器,确保共享订阅的消息能被正常接收 messageProducer.setMessageSelector(null) return IntegrationFlow.from(messageProducer) .channel("mqttInputChannel") .get() } @ServiceActivator(inputChannel = "mqttInputChannel") fun handler(message: Message<String>) { println("Received message: ${message.payload}") }
补充说明
如果需要更精细的消息过滤,可以自定义MessageSelector实现,基于实际的消息主题(如foo/test)进行匹配,替代默认的订阅主题匹配逻辑。
内容的提问来源于stack exchange,提问作者protenhan
相关产品推荐
相关产品推荐

