From afc02e72c178a3bf3433fce4fba73f16ae3f5b40 Mon Sep 17 00:00:00 2001 From: Maksym Dudnik Date: Wed, 13 Feb 2019 12:28:41 +0200 Subject: [PATCH 1/2] polygon color function fix for multiply datasources --- ui/src/app/widget/lib/map-widget2.js | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/ui/src/app/widget/lib/map-widget2.js b/ui/src/app/widget/lib/map-widget2.js index 4c6248f85b..7ac960447b 100644 --- a/ui/src/app/widget/lib/map-widget2.js +++ b/ui/src/app/widget/lib/map-widget2.js @@ -305,12 +305,9 @@ export default class TbMapWidgetV2 { } function updateLocationPolygonColor(location, color) { - if (!location.settings.calculatedPolygonColor || location.settings.calculatedPolygonColor !== color) { + if (location.polygon && color) { location.settings.calculatedPolygonColor = color; - if (location.polygon) { - tbMap.map.updatePolygonColor(location.polygon, location.settings, color); - } - + tbMap.map.updatePolygonColor(location.polygon, location.settings, color); } } @@ -338,10 +335,8 @@ export default class TbMapWidgetV2 { function updateLocationStyle(location, dataMap) { updateLocationLabel(location, dataMap); var color = calculateLocationColor(location, dataMap); - var polygonColor = calculateLocationPolygonColor(location, dataMap); var image = calculateLocationMarkerImage(location, dataMap); updateLocationColor(location, color, image); - if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, polygonColor); updateLocationMarkerIcon(location, image); } @@ -441,6 +436,7 @@ export default class TbMapWidgetV2 { if (location.marker) { updateLocationStyle(location, dataMap); } + } } return locationChanged; @@ -456,11 +452,13 @@ export default class TbMapWidgetV2 { locationPolygonClick(event, location); }, [location.dsIndex]); tbMap.polygons.push(location.polygon); + if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, calculateLocationPolygonColor(location, dataMap)); } else if (polygonLatLngs.length > 0) { let prevPolygonArr = tbMap.map.getPolygonLatLngs(location.polygon); if (!prevPolygonArr || !arraysEqual(prevPolygonArr, polygonLatLngs)) { tbMap.map.setPolygonLatLngs(location.polygon, polygonLatLngs); } + if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, calculateLocationPolygonColor(location, dataMap)); } } } From d863ecfa6b0996d55d75cc950ff56dd81842552a Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Wed, 13 Feb 2019 14:15:14 +0200 Subject: [PATCH 2/2] Improve rule node executors to use work stealing thread pool. Use shared event loop for Mqtt rule nodes. --- .../server/actors/ActorSystemContext.java | 5 ++ .../actors/ruleChain/DefaultTbContext.java | 6 +++ .../executors/AbstractListeningExecutor.java | 2 +- .../SharedEventLoopGroupService.java | 48 +++++++++++++++++++ .../AbstractNashornJsInvokeService.java | 2 +- .../src/main/resources/thingsboard.yml | 6 +-- .../rule/engine/api/TbContext.java | 4 ++ .../rule/engine/mqtt/TbMqttNode.java | 11 ++--- ui/package-lock.json | 2 +- 9 files changed, 72 insertions(+), 14 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 8433fdb1ea..2cc7e470c2 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -67,6 +67,7 @@ import org.thingsboard.server.service.encoding.DataDecodingEncodingService; import org.thingsboard.server.service.executors.ClusterRpcCallbackExecutorService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.ExternalCallExecutorService; +import org.thingsboard.server.service.executors.SharedEventLoopGroupService; import org.thingsboard.server.service.mail.MailExecutorService; import org.thingsboard.server.service.rpc.DeviceRpcService; import org.thingsboard.server.service.script.JsExecutorService; @@ -206,6 +207,10 @@ public class ActorSystemContext { @Getter private ExternalCallExecutorService externalCallExecutorService; + @Autowired + @Getter + private SharedEventLoopGroupService sharedEventLoopGroupService; + @Autowired @Getter private MailService mailService; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index ed50bdafa6..90a9d56e0c 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -17,6 +17,7 @@ package org.thingsboard.server.actors.ruleChain; import akka.actor.ActorRef; import com.datastax.driver.core.utils.UUIDs; +import io.netty.channel.EventLoopGroup; import org.springframework.util.StringUtils; import org.thingsboard.rule.engine.api.ListeningExecutor; import org.thingsboard.rule.engine.api.MailService; @@ -238,6 +239,11 @@ class DefaultTbContext implements TbContext { return mainCtx.getRuleChainTransactionService(); } + @Override + public EventLoopGroup getSharedEventLoop() { + return mainCtx.getSharedEventLoopGroupService().getSharedEventLoopGroup(); + } + @Override public MailService getMailService() { if (mainCtx.isAllowSystemMailService()) { diff --git a/application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java b/application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java index fabd345b64..221915d02c 100644 --- a/application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java +++ b/application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java @@ -34,7 +34,7 @@ public abstract class AbstractListeningExecutor implements ListeningExecutor { @PostConstruct public void init() { - this.service = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(getThreadPollSize())); + this.service = MoreExecutors.listeningDecorator(Executors.newWorkStealingPool(getThreadPollSize())); } @PreDestroy diff --git a/application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java b/application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java new file mode 100644 index 0000000000..e61b0db0f6 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java @@ -0,0 +1,48 @@ +/** + * Copyright © 2016-2019 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.executors; + +import com.google.common.util.concurrent.MoreExecutors; +import io.netty.channel.EventLoopGroup; +import io.netty.channel.nio.NioEventLoopGroup; +import lombok.Getter; +import org.springframework.stereotype.Component; + +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +@Component +public class SharedEventLoopGroupService { + + @Getter + private EventLoopGroup sharedEventLoopGroup; + + @PostConstruct + public void init() { + this.sharedEventLoopGroup = new NioEventLoopGroup(); + } + + @PreDestroy + public void destroy() { + if (this.sharedEventLoopGroup != null) { + this.sharedEventLoopGroup.shutdownGracefully(0, 5, TimeUnit.SECONDS); + } + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java index f365ef68dc..62e7c24b1a 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java @@ -42,7 +42,7 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer public void init() { if (useJsSandbox()) { sandbox = NashornSandboxes.create(); - monitorExecutorService = Executors.newFixedThreadPool(getMonitorThreadPoolSize()); + monitorExecutorService = Executors.newWorkStealingPool(getMonitorThreadPoolSize()); sandbox.setExecutor(monitorExecutorService); sandbox.setMaxCPUTime(getMaxCpuTime()); sandbox.allowNoBraces(false); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index ba9382aa2d..235646f6c2 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -200,13 +200,13 @@ actors: # Specify thread pool size for database request callbacks executor service db_callback_thread_pool_size: "${ACTORS_RULE_DB_CALLBACK_THREAD_POOL_SIZE:1}" # Specify thread pool size for javascript executor service - js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:10}" + js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:50}" # Specify thread pool size for mail sender executor service - mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:10}" + mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:50}" # Whether to allow usage of system mail service for rules allow_system_mail_service: "${ACTORS_RULE_ALLOW_SYSTEM_MAIL_SERVICE:true}" # Specify thread pool size for external call service - external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:10}" + external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:50}" chain: # Errors for particular actor are persisted once per specified amount of milliseconds error_persist_frequency: "${ACTORS_RULE_CHAIN_ERROR_FREQUENCY:3000}" diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index d3476a470c..9ebc610e1b 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.api; +import io.netty.channel.EventLoopGroup; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; @@ -105,4 +106,7 @@ public interface TbContext { String getNodeId(); RuleChainTransactionService getRuleChainTransactionService(); + + EventLoopGroup getSharedEventLoop(); + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java index 3ffa910802..dd17414b87 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java @@ -59,15 +59,13 @@ public class TbMqttNode implements TbNode { private TbMqttNodeConfiguration config; - private EventLoopGroup eventLoopGroup; private MqttClient mqttClient; @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { try { this.config = TbNodeUtils.convert(configuration, TbMqttNodeConfiguration.class); - this.eventLoopGroup = new NioEventLoopGroup(); - this.mqttClient = initClient(); + this.mqttClient = initClient(ctx); } catch (Exception e) { throw new TbNodeException(e); } @@ -99,12 +97,9 @@ public class TbMqttNode implements TbNode { if (this.mqttClient != null) { this.mqttClient.disconnect(); } - if (this.eventLoopGroup != null) { - this.eventLoopGroup.shutdownGracefully(0, 5, TimeUnit.SECONDS); - } } - private MqttClient initClient() throws Exception { + private MqttClient initClient(TbContext ctx) throws Exception { Optional sslContextOpt = initSslContext(); MqttClientConfig config = sslContextOpt.isPresent() ? new MqttClientConfig(sslContextOpt.get()) : new MqttClientConfig(); if (!StringUtils.isEmpty(this.config.getClientId())) { @@ -113,7 +108,7 @@ public class TbMqttNode implements TbNode { config.setCleanSession(this.config.isCleanSession()); this.config.getCredentials().configure(config); MqttClient client = MqttClient.create(config, null); - client.setEventLoop(this.eventLoopGroup); + client.setEventLoop(ctx.getSharedEventLoop()); Future connectFuture = client.connect(this.config.getHost(), this.config.getPort()); MqttConnectResult result; try { diff --git a/ui/package-lock.json b/ui/package-lock.json index 30fd603198..9b8224e0ff 100644 --- a/ui/package-lock.json +++ b/ui/package-lock.json @@ -1,6 +1,6 @@ { "name": "thingsboard", - "version": "2.3.0", + "version": "2.3.1", "lockfileVersion": 1, "requires": true, "dependencies": {