diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index c943e4c4d0..a5abd77ad5 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -52,6 +52,7 @@ import org.thingsboard.server.service.subscription.TbSubscriptionUtils; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -98,6 +99,11 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> pendingMap = msgs.stream().collect( Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); - ConcurrentMap> failedMap = new ConcurrentHashMap<>(); CountDownLatch processingTimeoutLatch = new CountDownLatch(1); + TbPackProcessingContext> ctx = new TbPackProcessingContext<>( + processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>()); pendingMap.forEach((id, msg) -> { log.trace("[{}] Creating main callback for message: {}", id, msg.getValue()); - TbCallback callback = new TbPackCallback<>(id, processingTimeoutLatch, pendingMap, failedMap); + TbCallback callback = new TbPackCallback<>(id, ctx); try { ToCoreMsg toCoreMsg = msg.getValue(); if (toCoreMsg.hasToSubscriptionMgrMsg()) { @@ -147,8 +154,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService log.warn("[{}] Timeout to process message: {}", id, msg.getValue())); - failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process message: {}", id, msg.getValue())); + ctx.getAckMap().forEach((id, msg) -> log.warn("[{}] Timeout to process message: {}", id, msg.getValue())); + ctx.getFailedMap().forEach((id, msg) -> log.warn("[{}] Failed to process message: {}", id, msg.getValue())); } mainConsumer.commit(); } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index 6736518bcc..0c48c54412 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -31,7 +31,6 @@ import org.thingsboard.server.common.msg.queue.ServiceQueue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbMsgCallback; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.queue.TbQueueConsumer; @@ -64,9 +63,6 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import java.util.function.BiConsumer; -import java.util.function.Function; -import java.util.stream.Collectors; @Service @TbRuleEngineComponent @@ -116,10 +112,10 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< @PreDestroy public void stop() { + super.destroy(); if (submitExecutor != null) { submitExecutor.shutdownNow(); } - ruleEngineSettings.getQueues().forEach(config -> consumerConfigurations.put(config.getName(), config)); } @@ -156,7 +152,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< submitStrategy.init(msgs); while (!stopped) { - ProcessingAttemptContext ctx = new ProcessingAttemptContext(submitStrategy); + TbMsgPackProcessingContext ctx = new TbMsgPackProcessingContext(submitStrategy); submitStrategy.submitAttempt((id, msg) -> submitExecutor.submit(() -> { log.trace("[{}] Creating callback for message: {}", id, msg.getValue()); ToRuleEngineMsg toRuleEngineMsg = msg.getValue(); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java index 2a6b6a658d..f093cc885a 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java @@ -18,21 +18,17 @@ package org.thingsboard.server.service.queue; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.RuleEngineException; -import org.thingsboard.server.common.msg.queue.RuleNodeException; -import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbMsgCallback; import java.util.UUID; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.CountDownLatch; @Slf4j public class TbMsgPackCallback implements TbMsgCallback { private final UUID id; private final TenantId tenantId; - private final ProcessingAttemptContext ctx; + private final TbMsgPackProcessingContext ctx; - public TbMsgPackCallback(UUID id, TenantId tenantId, ProcessingAttemptContext ctx) { + public TbMsgPackCallback(UUID id, TenantId tenantId, TbMsgPackProcessingContext ctx) { this.id = id; this.tenantId = tenantId; this.ctx = ctx; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ProcessingAttemptContext.java b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java similarity index 96% rename from application/src/main/java/org/thingsboard/server/service/queue/ProcessingAttemptContext.java rename to application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java index aefb1697cd..7e78ac6f5f 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ProcessingAttemptContext.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java @@ -29,7 +29,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; -public class ProcessingAttemptContext { +public class TbMsgPackProcessingContext { private final TbRuleEngineSubmitStrategy submitStrategy; @@ -44,7 +44,7 @@ public class ProcessingAttemptContext { @Getter private final ConcurrentMap exceptionsMap = new ConcurrentHashMap<>(); - public ProcessingAttemptContext(TbRuleEngineSubmitStrategy submitStrategy) { + public TbMsgPackProcessingContext(TbRuleEngineSubmitStrategy submitStrategy) { this.submitStrategy = submitStrategy; this.pendingMap = submitStrategy.getPendingMap(); this.pendingCount = new AtomicInteger(pendingMap.size()); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbPackCallback.java b/application/src/main/java/org/thingsboard/server/service/queue/TbPackCallback.java index ba5a883ea0..5a5c172ee6 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbPackCallback.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbPackCallback.java @@ -19,44 +19,26 @@ import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.msg.queue.TbCallback; import java.util.UUID; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.CountDownLatch; @Slf4j public class TbPackCallback implements TbCallback { - private final CountDownLatch processingTimeoutLatch; - private final ConcurrentMap ackMap; - private final ConcurrentMap failedMap; + private final TbPackProcessingContext ctx; private final UUID id; - public TbPackCallback(UUID id, - CountDownLatch processingTimeoutLatch, - ConcurrentMap ackMap, - ConcurrentMap failedMap) { + public TbPackCallback(UUID id, TbPackProcessingContext ctx) { this.id = id; - this.processingTimeoutLatch = processingTimeoutLatch; - this.ackMap = ackMap; - this.failedMap = failedMap; + this.ctx = ctx; } @Override public void onSuccess() { log.trace("[{}] ON SUCCESS", id); - T msg = ackMap.remove(id); - if (msg != null && ackMap.isEmpty()) { - processingTimeoutLatch.countDown(); - } + ctx.onSuccess(id); } @Override public void onFailure(Throwable t) { log.trace("[{}] ON FAILURE", id, t); - T msg = ackMap.remove(id); - if (msg != null) { - failedMap.put(id, msg); - } - if (ackMap.isEmpty()) { - processingTimeoutLatch.countDown(); - } + ctx.onFailure(id, t); } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbPackProcessingContext.java b/application/src/main/java/org/thingsboard/server/service/queue/TbPackProcessingContext.java new file mode 100644 index 0000000000..e9f1224625 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbPackProcessingContext.java @@ -0,0 +1,90 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.queue; + +import lombok.extern.slf4j.Slf4j; + +import java.util.UUID; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +@Slf4j +public class TbPackProcessingContext { + + private final AtomicInteger pendingCount; + private final CountDownLatch processingTimeoutLatch; + private final ConcurrentMap ackMap; + private final ConcurrentMap failedMap; + + public TbPackProcessingContext(CountDownLatch processingTimeoutLatch, + ConcurrentMap ackMap, + ConcurrentMap failedMap) { + this.processingTimeoutLatch = processingTimeoutLatch; + this.pendingCount = new AtomicInteger(ackMap.size()); + this.ackMap = ackMap; + this.failedMap = failedMap; + } + + public boolean await(long packProcessingTimeout, TimeUnit milliseconds) throws InterruptedException { + return processingTimeoutLatch.await(packProcessingTimeout, milliseconds); + } + + public void onSuccess(UUID id) { + boolean empty = false; + T msg = ackMap.remove(id); + if (msg != null) { + empty = pendingCount.decrementAndGet() == 0; + } + if (empty) { + processingTimeoutLatch.countDown(); + } else { + if (log.isTraceEnabled()) { + log.trace("Items left: {}", ackMap.size()); + for (T t : ackMap.values()) { + log.trace("left item: {}", t); + } + } + } + } + + public void onFailure(UUID id, Throwable t) { + boolean empty = false; + T msg = ackMap.remove(id); + if (msg != null) { + empty = pendingCount.decrementAndGet() == 0; + failedMap.put(id, msg); + if (log.isTraceEnabled()) { + log.trace("Items left: {}", ackMap.size()); + for (T v : ackMap.values()) { + log.trace("left item: {}", v); + } + } + } + if (empty) { + processingTimeoutLatch.countDown(); + } + } + + public ConcurrentMap getAckMap() { + return ackMap; + } + + public ConcurrentMap getFailedMap() { + return failedMap; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index a238e34337..c2705fcbdc 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -23,11 +23,13 @@ import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionChangeEvent; import org.thingsboard.server.service.encoding.DataDecodingEncodingService; import org.thingsboard.server.service.queue.TbPackCallback; +import org.thingsboard.server.service.queue.TbPackProcessingContext; import javax.annotation.PreDestroy; import java.util.List; @@ -92,11 +94,12 @@ public abstract class AbstractConsumerService> pendingMap = msgs.stream().collect( Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); - ConcurrentMap> failedMap = new ConcurrentHashMap<>(); CountDownLatch processingTimeoutLatch = new CountDownLatch(1); + TbPackProcessingContext> ctx = new TbPackProcessingContext<>( + processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>()); pendingMap.forEach((id, msg) -> { log.trace("[{}] Creating notification callback for message: {}", id, msg.getValue()); - TbCallback callback = new TbPackCallback<>(id, processingTimeoutLatch, pendingMap, failedMap); + TbCallback callback = new TbPackCallback<>(id, ctx); try { handleNotification(id, msg, callback); } catch (Throwable e) { @@ -105,8 +108,8 @@ public abstract class AbstractConsumerService log.warn("[{}] Timeout to process notification: {}", id, msg.getValue())); - failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process notification: {}", id, msg.getValue())); + ctx.getAckMap().forEach((id, msg) -> log.warn("[{}] Timeout to process notification: {}", id, msg.getValue())); + ctx.getFailedMap().forEach((id, msg) -> log.warn("[{}] Failed to process notification: {}", id, msg.getValue())); } nfConsumer.commit(); } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingResult.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingResult.java index 8e0fcaa74a..e818ffb9d7 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingResult.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingResult.java @@ -20,7 +20,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.service.queue.ProcessingAttemptContext; +import org.thingsboard.server.service.queue.TbMsgPackProcessingContext; import java.util.UUID; import java.util.concurrent.ConcurrentMap; @@ -32,9 +32,9 @@ public class TbRuleEngineProcessingResult { @Getter private final boolean timeout; @Getter - private final ProcessingAttemptContext ctx; + private final TbMsgPackProcessingContext ctx; - public TbRuleEngineProcessingResult(boolean timeout, ProcessingAttemptContext ctx) { + public TbRuleEngineProcessingResult(boolean timeout, TbMsgPackProcessingContext ctx) { this.timeout = timeout; this.ctx = ctx; this.success = !timeout && ctx.getPendingMap().isEmpty() && ctx.getFailedMap().isEmpty(); diff --git a/application/src/test/java/org/thingsboard/server/service/queue/ProcessingAttemptContextTest.java b/application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java similarity index 94% rename from application/src/test/java/org/thingsboard/server/service/queue/ProcessingAttemptContextTest.java rename to application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java index 2de2414f71..43cd62ea97 100644 --- a/application/src/test/java/org/thingsboard/server/service/queue/ProcessingAttemptContextTest.java +++ b/application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java @@ -37,7 +37,7 @@ import static org.mockito.Mockito.when; @Slf4j @RunWith(MockitoJUnitRunner.class) -public class ProcessingAttemptContextTest { +public class TbMsgPackProcessingContextTest { @Test public void testHighConcurrencyCase() throws InterruptedException { @@ -51,7 +51,7 @@ public class ProcessingAttemptContextTest { messages.put(UUID.randomUUID(), new TbProtoQueueMsg<>(UUID.randomUUID(), null)); } when(strategyMock.getPendingMap()).thenReturn(messages); - ProcessingAttemptContext context = new ProcessingAttemptContext(strategyMock); + TbMsgPackProcessingContext context = new TbMsgPackProcessingContext(strategyMock); for (UUID uuid : messages.keySet()) { for (int i = 0; i < parallelCount; i++) { executorService.submit(() -> context.onSuccess(uuid));