From 40bcd2fa8ab729898becb21b3286807cc6654057 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 30 Jul 2021 15:44:30 +0300 Subject: [PATCH] mqtt transport refactored msqProcessorExecutor lifecycle --- .../transport/mqtt/MqttTransportContext.java | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 64700b3923..8eae488be9 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -29,6 +29,7 @@ import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor; import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; import java.util.concurrent.ExecutorService; /** @@ -68,6 +69,22 @@ public class MqttTransportContext extends TransportContext { private int messageQueueSizePerDeviceLimit; @Getter - private final ExecutorService msqProcessorExecutor = ThingsBoardExecutors.newWorkStealingPool(Runtime.getRuntime().availableProcessors() + 1, "msg-processor-on-device-connect"); + private ExecutorService msqProcessorExecutor; + + @Override + @PostConstruct + public void init() { + super.init(); + msqProcessorExecutor = ThingsBoardExecutors.newWorkStealingPool(Runtime.getRuntime().availableProcessors() + 1, "msg-processor-on-device-connect"); + } + + @Override + @PreDestroy + public void stop() { + super.stop(); + if (msqProcessorExecutor != null) { + msqProcessorExecutor.shutdownNow(); + } + } }