diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java index 3f1c9cbcd0..86053614cc 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java @@ -19,6 +19,7 @@ package org.springframework.integration.mqtt.inbound; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.ConcurrentModificationException; import java.util.List; import java.util.Map; import java.util.concurrent.locks.Lock; @@ -289,7 +290,7 @@ public class Mqttv5PahoMessageDrivenChannelAdapter try { if (this.mqttClient != null && this.mqttClient.isConnected()) { if (this.connectionOptions.isCleanStart()) { - this.mqttClient.unsubscribe(topics).waitForCompletion(getCompletionTimeout()); + unsubscribe(topics); // Have to re-subscribe on next start if connection is not lost. this.readyToSubscribeOnStart = true; @@ -310,12 +311,22 @@ public class Mqttv5PahoMessageDrivenChannelAdapter } } + private void unsubscribe(String... topics) throws MqttException { + try { + // Catch ConcurrentModificationException: https://github.com/eclipse/paho.mqtt.java/issues/986 + this.mqttClient.unsubscribe(topics).waitForCompletion(getCompletionTimeout()); + } + catch (ConcurrentModificationException ex) { + logger.error(ex, () -> "Error unsubscribing from " + Arrays.toString(topics)); + } + } + @Override public void destroy() { super.destroy(); try { if (getClientManager() == null && this.mqttClient != null) { - this.mqttClient.close(true); + this.mqttClient.close(); } } catch (MqttException ex) { @@ -360,7 +371,7 @@ public class Mqttv5PahoMessageDrivenChannelAdapter this.topicLock.lock(); try { if (this.mqttClient != null && this.mqttClient.isConnected()) { - this.mqttClient.unsubscribe(topic).waitForCompletion(getCompletionTimeout()); + unsubscribe(topic); } super.removeTopic(topic); if (!CollectionUtils.isEmpty(this.subscriptions)) {