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
相关产品推荐
相关产品推荐

