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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:47:50