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 084b43f4b6..2c779d5d5a 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 @@ -85,59 +85,75 @@ public class DefaultTbQueueRequestTemplate { - long nextCleanupMs = 0L; - while (!stopped) { - try { - final int pendingRequestsCount = pendingRequests.size(); - log.trace("Starting template pool topic {}, for pendingRequests {}", responseTemplate.getTopic(), pendingRequestsCount); - List responses = responseTemplate.poll(pollInterval); //poll js responses - log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); - responses.forEach(response -> { - byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER); - UUID requestId; - if (requestIdHeader == null) { - log.error("[{}] Missing requestId in header and body", response); - } else { - requestId = bytesToUuid(requestIdHeader); - log.trace("[{}] Response received: {}", requestId, response); - ResponseMetaData expectedResponse = pendingRequests.remove(requestId); - if (expectedResponse == null) { - log.trace("[{}] Invalid or stale request", requestId); - } else { - expectedResponse.future.set(response); + executor.submit(this::fetchAndProcessResponses); + } + + void fetchAndProcessResponses() { + long nextCleanupMs = 0L; + while (!stopped) { + try { + final int pendingRequestsCount = pendingRequests.size(); + log.trace("Starting template pool topic {}, for pendingRequests {}", responseTemplate.getTopic(), pendingRequestsCount); + List responses = doPoll(); //poll js responses + log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); + responses.forEach(this::processResponse); + responseTemplate.commit(); + tickTs = System.currentTimeMillis(); + tickSize = pendingRequests.size(); + if (nextCleanupMs < tickTs) { + //cleanup; + pendingRequests.forEach((key, value) -> { + if (value.expTime < tickTs) { + ResponseMetaData staleRequest = pendingRequests.remove(key); + if (staleRequest != null) { + setTimeoutException(key, staleRequest); } } }); - responseTemplate.commit(); - tickTs = System.currentTimeMillis(); - tickSize = pendingRequests.size(); - if (nextCleanupMs < tickTs) { - //cleanup; - pendingRequests.forEach((key, value) -> { - if (value.expTime < tickTs) { - ResponseMetaData staleRequest = pendingRequests.remove(key); - if (staleRequest != null) { - log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); - staleRequest.future.setException(new TimeoutException()); - } - } - }); - nextCleanupMs = tickTs + maxRequestTimeout; - } - } catch (Throwable e) { - log.warn("Failed to obtain responses from queue. Going to sleep " + pollInterval + "ms", e); - try { - Thread.sleep(pollInterval); - } catch (InterruptedException e2) { - log.trace("Failed to wait until the server has capacity to handle new responses", e2); - } + nextCleanupMs = tickTs + maxRequestTimeout; } + } catch (Throwable e) { + log.warn("Failed to obtain responses from queue. Going to sleep " + pollInterval + "ms", e); + sleep(); } - }); + } + } + + List doPoll() { + return responseTemplate.poll(pollInterval); + } + + void sleep() { + try { + Thread.sleep(pollInterval); + } catch (InterruptedException e2) { + log.trace("Failed to wait until the server has capacity to handle new responses", e2); + } + } + + void setTimeoutException(UUID key, ResponseMetaData staleRequest) { + log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); + staleRequest.future.setException(new TimeoutException()); + } + + void processResponse(Response response) { + byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER); + UUID requestId; + if (requestIdHeader == null) { + log.error("[{}] Missing requestId in header and body", response); + } else { + requestId = bytesToUuid(requestIdHeader); + log.trace("[{}] Response received: {}", requestId, response); + ResponseMetaData expectedResponse = pendingRequests.remove(requestId); + if (expectedResponse == null) { + log.warn("[{}] Invalid or stale request, response: {}", requestId, response); + } else { + expectedResponse.future.set(response); + } + } } @Override @@ -174,7 +190,7 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); ResponseMetaData responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future); pendingRequests.putIfAbsent(requestId, responseMetaData); - log.trace("[{}] Sending request, key [{}], expTime [{}]", requestId, request.getKey(), responseMetaData.expTime); + log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, request); if (messagesStats != null) { messagesStats.incrementTotal(); } @@ -184,7 +200,7 @@ public class DefaultTbQueueRequestTemplate