From 0d94b6231e8fd19d0673db256534d4eaa95572c3 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 15 Jun 2021 15:22:30 +0300 Subject: [PATCH] rule-engine - fixed name for reusable threads --- .../queue/DefaultTbRuleEngineConsumerService.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) 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 4684db3e2e..a4395f0d3c 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 @@ -82,6 +82,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< public static final String SUCCESSFUL_STATUS = "successful"; public static final String FAILED_STATUS = "failed"; + public static final String THREAD_TOPIC_SPLITERATOR = " | "; @Value("${queue.rule-engine.poll-interval}") private long pollDuration; @Value("${queue.rule-engine.pack-processing-timeout}") @@ -245,7 +246,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< } void consumerLoop(TbQueueConsumer> consumer, TbRuleEngineQueueConfiguration configuration, TbRuleEngineConsumerStats stats, String threadSuffix) { - Thread.currentThread().setName("" + Thread.currentThread().getName() + "-" + threadSuffix); + updateCurrentThreadName(threadSuffix); while (!stopped && !consumer.isStopped()) { try { List> msgs = consumer.poll(pollDuration); @@ -299,6 +300,16 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< log.info("TB Rule Engine Consumer stopped."); } + void updateCurrentThreadName(String threadSuffix) { + String name = Thread.currentThread().getName(); + int spliteratorIndex = name.indexOf(THREAD_TOPIC_SPLITERATOR); + if (spliteratorIndex > 0) { + name = name.substring(0, spliteratorIndex); + } + name = name + THREAD_TOPIC_SPLITERATOR + threadSuffix; + Thread.currentThread().setName(name); + } + TbRuleEngineProcessingStrategy getAckStrategy(TbRuleEngineQueueConfiguration configuration) { return processingStrategyFactory.newInstance(configuration.getName(), configuration.getProcessingStrategy()); }