From e33709d830d25d8fa9b998e80f4898168db63b63 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Thu, 27 Aug 2020 17:11:00 +0300 Subject: [PATCH 1/3] fix/getAttributesRequest & fix/hsql-inserts-with-json (#3376) * fix GetAttributeRequest & Hsql inserts * added deleted keys to GetAttributeResponseMsg for gateway * changed logic for handleGetAttributesRequest on deleted attributes * removed unused classes and remove deleted keys in GetAttributeResponseMsg proto * revert yml file * removed getDeletedAttributeKeysCount method call in Coap transport --- .../device/DeviceActorMessageProcessor.java | 21 ++++---- .../transport/DefaultTransportApiService.java | 2 +- .../server/common/msg/kv/AttributesKVMsg.java | 29 ----------- .../common/msg/kv/BasicAttributeKVMsg.java | 52 ------------------- common/queue/src/main/proto/queue.proto | 1 - .../coap/adaptors/JsonCoapAdaptor.java | 2 +- .../transport/adaptor/JsonConverter.java | 35 ------------- .../dao/sqlts/hsql/JpaHsqlTimeseriesDao.java | 1 + .../hsql/HsqlLatestInsertTsRepository.java | 4 +- .../resources/sql/hsql/drop-all-tables.sql | 1 + 10 files changed, 16 insertions(+), 132 deletions(-) delete mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/kv/AttributesKVMsg.java delete mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/kv/BasicAttributeKVMsg.java diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index e328931c0e..cf4aa3edda 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -28,6 +28,7 @@ import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActorCtx; import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; @@ -79,8 +80,6 @@ import java.util.UUID; import java.util.function.Consumer; import java.util.stream.Collectors; -import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE; -import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE; /** * @author Andrew Shvayka @@ -279,17 +278,17 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { ListenableFuture> clientAttributesFuture; ListenableFuture> sharedAttributesFuture; if (CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) { - clientAttributesFuture = findAllAttributesByScope(CLIENT_SCOPE); - sharedAttributesFuture = findAllAttributesByScope(SHARED_SCOPE); + clientAttributesFuture = findAllAttributesByScope(DataConstants.CLIENT_SCOPE); + sharedAttributesFuture = findAllAttributesByScope(DataConstants.SHARED_SCOPE); } else if (!CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && !CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) { - clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), CLIENT_SCOPE); - sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), SHARED_SCOPE); + clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), DataConstants.CLIENT_SCOPE); + sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), DataConstants.SHARED_SCOPE); } else if (CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && !CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) { clientAttributesFuture = Futures.immediateFuture(Collections.emptyList()); - sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), SHARED_SCOPE); + sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), DataConstants.SHARED_SCOPE); } else { sharedAttributesFuture = Futures.immediateFuture(Collections.emptyList()); - clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), CLIENT_SCOPE); + clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), DataConstants.CLIENT_SCOPE); } return Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture)); } @@ -316,7 +315,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder(); if (msg.isDeleted()) { List sharedKeys = msg.getDeletedKeys().stream() - .filter(key -> SHARED_SCOPE.equals(key.getScope())) + .filter(key -> DataConstants.SHARED_SCOPE.equals(key.getScope())) .map(AttributeKey::getAttributeKey) .collect(Collectors.toList()); if (!sharedKeys.isEmpty()) { @@ -324,7 +323,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { hasNotificationData = true; } } else { - if (SHARED_SCOPE.equals(msg.getScope())) { + if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) { List attributes = new ArrayList<>(msg.getValues()); if (attributes.size() > 0) { List sharedUpdated = msg.getValues().stream().map(this::toTsKvProto) @@ -334,7 +333,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { hasNotificationData = true; } } else { - log.debug("[{}] No public server side attributes changed!", deviceId); + log.debug("[{}] No public shared side attributes changed!", deviceId); } } } diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index b08ff7c20f..3aa14ed639 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -154,7 +154,7 @@ public class DefaultTransportApiService implements TransportApiService { return TransportApiResponseMsg.newBuilder() .setGetOrCreateDeviceResponseMsg(GetOrCreateDeviceFromGatewayResponseMsg.newBuilder().setDeviceInfo(getDeviceInfoProto(device)).build()).build(); } catch (JsonProcessingException e) { - log.warn("[{}] Failed to lookup device by gateway id and name", gatewayId, requestMsg.getDeviceName(), e); + log.warn("[{}][{}] Failed to lookup device by gateway id and name", gatewayId, requestMsg.getDeviceName(), e); throw new RuntimeException(e); } finally { deviceCreationLock.unlock(); diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/kv/AttributesKVMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/kv/AttributesKVMsg.java deleted file mode 100644 index 85283242f9..0000000000 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/kv/AttributesKVMsg.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * 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.common.msg.kv; - -import java.io.Serializable; -import java.util.List; - -import org.thingsboard.server.common.data.kv.AttributeKey; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; - -public interface AttributesKVMsg extends Serializable { - - List getClientAttributes(); - List getSharedAttributes(); - List getDeletedAttributes(); -} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/kv/BasicAttributeKVMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/kv/BasicAttributeKVMsg.java deleted file mode 100644 index 8eb087f961..0000000000 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/kv/BasicAttributeKVMsg.java +++ /dev/null @@ -1,52 +0,0 @@ -/** - * 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.common.msg.kv; - -import lombok.AccessLevel; -import lombok.Data; -import lombok.RequiredArgsConstructor; -import org.thingsboard.server.common.data.kv.AttributeKey; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; - -import java.util.Collections; -import java.util.List; - -@Data -@RequiredArgsConstructor(access = AccessLevel.PRIVATE) -public class BasicAttributeKVMsg implements AttributesKVMsg { - - private static final long serialVersionUID = 1L; - - private final List clientAttributes; - private final List sharedAttributes; - private final List deletedAttributes; - - public static BasicAttributeKVMsg fromClient(List attributes) { - return new BasicAttributeKVMsg(attributes, Collections.emptyList(), Collections.emptyList()); - } - - public static BasicAttributeKVMsg fromShared(List attributes) { - return new BasicAttributeKVMsg(Collections.emptyList(), attributes, Collections.emptyList()); - } - - public static BasicAttributeKVMsg from(List client, List shared) { - return new BasicAttributeKVMsg(client, shared, Collections.emptyList()); - } - - public static AttributesKVMsg fromDeleted(List shared) { - return new BasicAttributeKVMsg(Collections.emptyList(), Collections.emptyList(), shared); - } -} diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index d3e850b39c..59d45d7db6 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/common/queue/src/main/proto/queue.proto @@ -127,7 +127,6 @@ message GetAttributeResponseMsg { int32 requestId = 1; repeated TsKvProto clientAttributeList = 2; repeated TsKvProto sharedAttributeList = 3; - repeated string deletedAttributeKeys = 4; string error = 5; } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java index 9d509d8f60..5c9c471570 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java @@ -125,7 +125,7 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { @Override public Response convertToPublish(CoapTransportResource.CoapSessionListener session, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException { - if (msg.getClientAttributeListCount() == 0 && msg.getSharedAttributeListCount() == 0 && msg.getDeletedAttributeKeysCount() == 0) { + if (msg.getClientAttributeListCount() == 0 && msg.getSharedAttributeListCount() == 0) { return new Response(CoAP.ResponseCode.NOT_FOUND); } else { Response response = new Response(CoAP.ResponseCode.CONTENT); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java index 8375b84ffa..a43437bb9d 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java @@ -35,7 +35,6 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.msg.kv.AttributesKVMsg; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; @@ -269,11 +268,6 @@ public class JsonConverter { payload.getSharedAttributeListList().forEach(addToObjectFromProto(attrObject)); result.add("shared", attrObject); } - if (payload.getDeletedAttributeKeysCount() > 0) { - JsonArray attrObject = new JsonArray(); - payload.getDeletedAttributeKeysList().forEach(attrObject::add); - result.add("deleted", attrObject); - } return result; } @@ -290,31 +284,6 @@ public class JsonConverter { return result; } - public static JsonObject toJson(AttributesKVMsg payload, boolean asMap) { - JsonObject result = new JsonObject(); - if (asMap) { - if (!payload.getClientAttributes().isEmpty()) { - JsonObject attrObject = new JsonObject(); - payload.getClientAttributes().forEach(addToObject(attrObject)); - result.add("client", attrObject); - } - if (!payload.getSharedAttributes().isEmpty()) { - JsonObject attrObject = new JsonObject(); - payload.getSharedAttributes().forEach(addToObject(attrObject)); - result.add("shared", attrObject); - } - } else { - payload.getClientAttributes().forEach(addToObject(result)); - payload.getSharedAttributes().forEach(addToObject(result)); - } - if (!payload.getDeletedAttributes().isEmpty()) { - JsonArray attrObject = new JsonArray(); - payload.getDeletedAttributes().forEach(addToObject(attrObject)); - result.add("deleted", attrObject); - } - return result; - } - public static JsonObject getJsonObjectForGateway(String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) { JsonObject result = new JsonObject(); result.addProperty("id", responseMsg.getRequestId()); @@ -370,10 +339,6 @@ public class JsonConverter { } } - private static Consumer addToObject(JsonArray result) { - return key -> result.add(key.getAttributeKey()); - } - private static Consumer addToObjectFromProto(JsonObject result) { return de -> { switch (de.getKv().getType()) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java index 5d94d12a1a..5ea387328d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java @@ -46,6 +46,7 @@ public class JpaHsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDa entity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null)); entity.setLongValue(tsKvEntry.getLongValue().orElse(null)); entity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null)); + entity.setJsonValue(tsKvEntry.getJsonValue().orElse(null)); log.trace("Saving entity: {}", entity); return tsQueue.add(entity); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/hsql/HsqlLatestInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/hsql/HsqlLatestInsertTsRepository.java index 224dc52805..429d5de31b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/hsql/HsqlLatestInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/hsql/HsqlLatestInsertTsRepository.java @@ -41,8 +41,8 @@ public class HsqlLatestInsertTsRepository extends AbstractInsertRepository imple "ON (ts_kv_latest.entity_id=T.entity_id " + "AND ts_kv_latest.key=T.key) " + "WHEN MATCHED THEN UPDATE SET ts_kv_latest.ts = T.ts, ts_kv_latest.bool_v = T.bool_v, ts_kv_latest.str_v = T.str_v, ts_kv_latest.long_v = T.long_v, ts_kv_latest.dbl_v = T.dbl_v, ts_kv_latest.json_v = T.json_v " + - "WHEN NOT MATCHED THEN INSERT (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " + - "VALUES (T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v);"; + "WHEN NOT MATCHED THEN INSERT (entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) " + + "VALUES (T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v, T.json_v);"; @Override public void saveOrUpdate(List entities) { diff --git a/dao/src/test/resources/sql/hsql/drop-all-tables.sql b/dao/src/test/resources/sql/hsql/drop-all-tables.sql index 1bdc1a7ece..778a4bbd9a 100644 --- a/dao/src/test/resources/sql/hsql/drop-all-tables.sql +++ b/dao/src/test/resources/sql/hsql/drop-all-tables.sql @@ -13,6 +13,7 @@ DROP TABLE IF EXISTS relation; DROP TABLE IF EXISTS tb_user; DROP TABLE IF EXISTS tenant; DROP TABLE IF EXISTS ts_kv; +DROP TABLE IF EXISTS ts_kv_dictionary; DROP TABLE IF EXISTS ts_kv_latest; DROP TABLE IF EXISTS user_credentials; DROP TABLE IF EXISTS widget_type; From 37940867263ed7cc8bca68187a68e3401460b815 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 28 Aug 2020 15:24:31 +0300 Subject: [PATCH 2/3] added logs for in memory queue --- .../server/queue/memory/InMemoryStorage.java | 18 ++++++++++++++++++ .../queue/memory/InMemoryTbQueueProducer.java | 2 +- 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java index f51f5b8ae0..c4414089c4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java @@ -23,16 +23,29 @@ import java.util.Collections; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @Slf4j public final class InMemoryStorage { private static InMemoryStorage instance; private final ConcurrentHashMap> storage; + private static ScheduledExecutorService statExecutor; private InMemoryStorage() { storage = new ConcurrentHashMap<>(); + statExecutor = Executors.newSingleThreadScheduledExecutor(); + statExecutor.scheduleAtFixedRate(this::printStats, 30, 30, TimeUnit.SECONDS); + } + + private void printStats() { + storage.forEach((topic, queue) -> { + if (queue.size() > 0) { + log.debug("Topic: [{}], Queue size: [{}]", topic, queue.size()); + } + }); } public static InMemoryStorage getInstance() { @@ -77,4 +90,9 @@ public final class InMemoryStorage { storage.clear(); } + public void destroy() { + if (statExecutor != null) { + statExecutor.shutdownNow(); + } + } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java index cfcd788a16..84a9a1fdf0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java @@ -53,6 +53,6 @@ public class InMemoryTbQueueProducer implements TbQueuePro @Override public void stop() { - + storage.destroy(); } } From 17ee07e09a0ae11b18e9aeb280cbd7a79afe3e8b Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 28 Aug 2020 16:01:43 +0300 Subject: [PATCH 3/3] changed log time --- .../org/thingsboard/server/queue/memory/InMemoryStorage.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java index c4414089c4..994ca26305 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java @@ -37,7 +37,7 @@ public final class InMemoryStorage { private InMemoryStorage() { storage = new ConcurrentHashMap<>(); statExecutor = Executors.newSingleThreadScheduledExecutor(); - statExecutor.scheduleAtFixedRate(this::printStats, 30, 30, TimeUnit.SECONDS); + statExecutor.scheduleAtFixedRate(this::printStats, 60, 60, TimeUnit.SECONDS); } private void printStats() {