From 1e066f21569aa7824177a6080577ac83d29ef2a7 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 12:41:50 +0300 Subject: [PATCH] added timeout parameter to the TbQueueRequestTemplate.send( ) --- .../server/queue/TbQueueRequestTemplate.java | 2 ++ .../common/DefaultTbQueueRequestTemplate.java | 18 +++++++++++++++--- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java index 5dc89a9c26..192f8e1675 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java @@ -24,6 +24,8 @@ public interface TbQueueRequestTemplate send(Request request); + ListenableFuture send(Request request, long timeoutNs); + void stop(); void setMessagesStats(MessagesStats messagesStats); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index c31dc3a899..5b03e1bc2d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -206,6 +206,11 @@ public class DefaultTbQueueRequestTemplate send(Request request) { + return send(request, this.maxRequestTimeoutNs); + } + + @Override + public ListenableFuture send(Request request, long requestTimeoutNs) { if (pendingRequests.mappingCount() >= maxPendingRequests) { log.warn("Pending request map is full [{}]! Consider to increase maxPendingRequests or increase processing performance", maxPendingRequests); return Futures.immediateFailedFuture(new RuntimeException("Pending request map is full!")); @@ -216,7 +221,7 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); - ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + maxRequestTimeoutNs, future, currentClockNs, maxRequestTimeoutNs); + ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + requestTimeoutNs, future, currentClockNs, requestTimeoutNs); log.trace("pending {}", responseMetaData); if (pendingRequests.putIfAbsent(requestId, responseMetaData) != null) { log.warn("Pending request already exists [{}]!", maxPendingRequests); @@ -226,11 +231,18 @@ public class DefaultTbQueueRequestTemplate