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

Spring Integration MQTT v5共享订阅无法正常工作问题求助

MQTT v5共享订阅无法接收消息问题排查与解决

问题描述

通过Spring Integration实现MQTT v5共享订阅,订阅主题为$share/group1/test,但共享订阅无法正常接收消息,仅非共享订阅可正常工作。

代码实现

@Bean
public ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager() {
    MqttConnectionOptions connectionOptions = new MqttConnectionOptions();
    connectionOptions.setServerURIs(new String[]{mqttProperties.getUrl()});
    connectionOptions.setConnectionTimeout(30);
    connectionOptions.setMaxReconnectDelay(30);
    connectionOptions.setUserName(mqttProperties.getUsername());
    connectionOptions.setPassword(mqttProperties
            .getPassword()
            .getBytes(StandardCharsets.UTF_8)
    );
    connectionOptions.setAutomaticReconnect(true);
    Mqttv5ClientManager clientManager =
            new Mqttv5ClientManager(
                    connectionOptions,
                    Objects.requireNonNull(MyUtils.getClientId()));
    clientManager.setPersistence(new MemoryPersistence());
    return clientManager;
}


@Bean
public IntegrationFlow mqttInFlow(
        final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) {
    Mqttv5PahoMessageDrivenChannelAdapter messageProducer =
            new Mqttv5PahoMessageDrivenChannelAdapter(
                    clientManager, "$share/group1/test");

    messageProducer.setPayloadType(byte[].class);
    messageProducer.setManualAcks(true);
    messageProducer.setMessageConverter(new ByteArrayMessageConverter());

    return IntegrationFlow.from(messageProducer)
            .handle(mqttMessageHandler)
            .get();
}


public class MqttMessageHandler implements GenericHandler<byte[]> {

    @Override
    public Object handle(final byte[] payload, final MessageHeaders headers) {
        log.info("Received message: {}", payload);
        
    }
}

已知条件

代码中订阅的是共享主题$share/group1/test

预期效果

发布到test主题的消息应能被MqttMessageHandler的handle方法接收并处理

问题现状

共享订阅无法正常接收消息,但非共享订阅可正常工作


问题分析与解决

核心原因

直接将$share/group1/test作为主题字符串传入Mqttv5PahoMessageDrivenChannelAdapter时,Spring Integration会对主题进行额外处理,导致实际发送给MQTT Broker的订阅格式不符合MQTT v5共享订阅的规范要求。

解决方案

通过Mqttv5Subscription对象配置共享订阅,而非手动拼接主题字符串,修改mqttInFlow实现:

@Bean
public IntegrationFlow mqttInFlow(
        final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) {
    // 创建订阅配置,指定原始主题
    Mqttv5Subscription subscription = new Mqttv5Subscription("test");
    // 设置共享订阅组名
    subscription.setSharedSubscriptionName("group1");
    // 根据实际需求设置QoS级别
    subscription.setQos(1);

    Mqttv5PahoMessageDrivenChannelAdapter messageProducer =
            new Mqttv5PahoMessageDrivenChannelAdapter(clientManager);
    // 添加配置好的共享订阅
    messageProducer.addSubscription(subscription);

    messageProducer.setPayloadType(byte[].class);
    messageProducer.setManualAcks(true);
    messageProducer.setMessageConverter(new ByteArrayMessageConverter());

    return IntegrationFlow.from(messageProducer)
            .handle(mqttMessageHandler)
            .get();
}

补充说明

  1. 使用Mqttv5Subscription并设置sharedSubscriptionName后,Spring Integration会自动按照MQTT v5规范构造正确的共享订阅主题格式($share/group1/test),无需手动拼接。
  2. 需确保使用的MQTT Broker支持MQTT v5协议及共享订阅功能(如EMQX、Mosquitto 2.0+等)。
  3. 若开启了手动确认(setManualAcks(true)),需在MqttMessageHandler的handle方法中手动确认消息,避免消息堆积或重复:
@Override
public Object handle(final byte[] payload, final MessageHeaders headers) {
    log.info("Received message: {}", new String(payload, StandardCharsets.UTF_8));
    // 手动确认消息
    MqttMessageAcknowledgment acknowledgment = headers.get(MqttHeaders.MESSAGE_ACKNOWLEDGMENT, MqttMessageAcknowledgment.class);
    if (acknowledgment != null) {
        acknowledgment.acknowledge();
    }
    return null;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 10:12:29