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

Spring Integration中如何自定义MqttSubscription设置no-local属性?

问题:Spring Integration中为MQTTv5订阅设置no-local选项

我在Spring Integration中使用org.eclipse.paho.mqttv5.client,尝试设置MQTT的no-local选项,代码如下:

@Bean
public MessageProducer inbound(ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) {
    Mqttv5PahoMessageDrivenChannelAdapter adapter = new Mqttv5PahoMessageDrivenChannelAdapter(
            clientManager,
            "test"
    );
    adapter.setCompletionTimeout(5000);
    adapter.setQos(2);
    adapter.connectComplete(true);
    adapter.setOutputChannel(mqttInputChannel());
    return adapter;
}

但Mqttv5PahoMessageDrivenChannelAdapter没有直接设置MqttSubscription(包含no-local配置)的方法。查看该类的subscribe方法源码发现,它仅用topic和qos创建订阅:

private void subscribe() {
    var clientManager = getClientManager();
    if (clientManager != null && this.mqttClient == null) {
        this.mqttClient = clientManager.getClient();
    }

    String[] topics = getTopic();
    ApplicationEventPublisher applicationEventPublisher = getApplicationEventPublisher();
    this.topicLock.lock();
    try {
        if (topics.length == 0) {
            return;
        }

        int[] requestedQos = getQos();
        MqttSubscription[] subscriptions = IntStream.range(0, topics.length)
                .mapToObj(i -> new MqttSubscription(topics[i], requestedQos[i]))
                .toArray(MqttSubscription[]::new);
        IMqttMessageListener listener = this::messageArrived;
        IMqttMessageListener[] listeners = IntStream.range(0, topics.length)
                .mapToObj(t -> listener)
                .toArray(IMqttMessageListener[]::new);
        this.mqttClient.subscribe(subscriptions, null, null, listeners, null)
                .waitForCompletion(getCompletionTimeout());
        String message = "Connected and subscribed to " + Arrays.toString(topics);
        logger.debug(message);
        if (applicationEventPublisher != null) {
            applicationEventPublisher.publishEvent(new MqttSubscribedEvent(this, message));
        }
    }
    catch (MqttException ex) {
        if (applicationEventPublisher != null) {
            applicationEventPublisher.publishEvent(new MqttConnectionFailedEvent(this, ex));
        }
        logger.error(ex, () -> "Error subscribing to " + Arrays.toString(topics));
    }
    finally {
        this.topicLock.unlock();
    }
}

解决方案

要设置no-local选项,需要自定义Mqttv5PahoMessageDrivenChannelAdapter的子类,重写subscribe方法,在创建MqttSubscription时指定no-local参数:

步骤1:创建自定义适配器子类

public class CustomMqttv5PahoMessageDrivenChannelAdapter extends Mqttv5PahoMessageDrivenChannelAdapter {

    private boolean noLocal;

    public CustomMqttv5PahoMessageDrivenChannelAdapter(ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager, String... topics) {
        super(clientManager, topics);
    }

    public void setNoLocal(boolean noLocal) {
        this.noLocal = noLocal;
    }

    @Override
    protected void subscribe() {
        var clientManager = getClientManager();
        if (clientManager != null && this.mqttClient == null) {
            this.mqttClient = clientManager.getClient();
        }

        String[] topics = getTopic();
        ApplicationEventPublisher applicationEventPublisher = getApplicationEventPublisher();
        this.topicLock.lock();
        try {
            if (topics.length == 0) {
                return;
            }

            int[] requestedQos = getQos();
            // 创建订阅时设置no-local选项
            MqttSubscription[] subscriptions = IntStream.range(0, topics.length)
                    .mapToObj(i -> {
                        MqttSubscription subscription = new MqttSubscription(topics[i], requestedQos[i]);
                        subscription.setNoLocal(this.noLocal);
                        return subscription;
                    })
                    .toArray(MqttSubscription[]::new);
            IMqttMessageListener listener = this::messageArrived;
            IMqttMessageListener[] listeners = IntStream.range(0, topics.length)
                    .mapToObj(t -> listener)
                    .toArray(IMqttMessageListener[]::new);
            this.mqttClient.subscribe(subscriptions, null, null, listeners, null)
                    .waitForCompletion(getCompletionTimeout());
            String message = "Connected and subscribed to " + Arrays.toString(topics);
            logger.debug(message);
            if (applicationEventPublisher != null) {
                applicationEventPublisher.publishEvent(new MqttSubscribedEvent(this, message));
            }
        } catch (MqttException ex) {
            if (applicationEventPublisher != null) {
                applicationEventPublisher.publishEvent(new MqttConnectionFailedEvent(this, ex));
            }
            logger.error(ex, () -> "Error subscribing to " + Arrays.toString(topics));
        } finally {
            this.topicLock.unlock();
        }
    }
}

步骤2:使用自定义适配器配置Bean

替换原有的Mqttv5PahoMessageDrivenChannelAdapter为自定义子类,并调用setNoLocal(true):

@Bean
public MessageProducer inbound(ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) {
    CustomMqttv5PahoMessageDrivenChannelAdapter adapter = new CustomMqttv5PahoMessageDrivenChannelAdapter(
            clientManager,
            "test"
    );
    adapter.setCompletionTimeout(5000);
    adapter.setQos(2);
    adapter.connectComplete(true);
    adapter.setNoLocal(true); // 设置no-local选项
    adapter.setOutputChannel(mqttInputChannel());
    return adapter;
}

说明

  • MqttSubscription的setNoLocal(boolean)方法用于控制是否接收当前客户端发布的消息:设置为true时,客户端不会收到自己发布的消息;默认是false。
  • 重写subscribe方法时,保留原有逻辑,仅修改MqttSubscription的创建逻辑,添加no-local配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 15:07:15