From 0b07c120256740e0f8ecc8f441ab81ef313b229e Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 26 Jun 2024 14:23:53 -0400 Subject: [PATCH] GH-9276: Mitigate `ConcurrentModificationException` in the `Mqttv5PahoMessageDrivenChannelAdapter` Fixes: #9276 The current Eclipse Paho client has wrong removal from map logic which leads to the `ConcurrentModificationException` when we unsubscribe from topic. * Catch a `ConcurrentModificationException` `this.mqttClient.unsubscribe()` and just log it as an `error` (cherry picked from commit 12fc353db13db5e4b8a69424520c772704634853) # Conflicts: # spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java --- .../Mqttv5PahoMessageDrivenChannelAdapter.java | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) 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 2fd8da914c..6599c4e649 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 @@ -17,6 +17,7 @@ package org.springframework.integration.mqtt.inbound; import java.util.Arrays; +import java.util.ConcurrentModificationException; import java.util.List; import java.util.Map; import java.util.concurrent.locks.Lock; @@ -240,7 +241,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; @@ -261,12 +262,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) { @@ -301,7 +312,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); }