From 03f5375a02acb3bcca7039e872165768529c47db Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Fri, 14 Feb 2020 19:18:18 +0200 Subject: [PATCH] JSON support (#2415) * Created JsonDataEntry and added DataType JSON * Added json to ts and attributes, created sql schema-entities-hsql.sql (json_v varchar) * refactored * refactored * added json array support * Aggregation improvement * Changed in JsonDataEntry value type from JsonNode to String * fix AggregatePartitionsFunction Co-authored-by: Yevhen Bondarenko <56396344+YevhenBondarenko@users.noreply.github.com> --- .../device/DeviceActorMessageProcessor.java | 7 + .../server/controller/BaseController.java | 2 + .../controller/TelemetryController.java | 43 ++- .../DefaultTelemetrySubscriptionService.java | 26 +- application/src/main/proto/cluster.proto | 1 + .../controller/ControllerSqlTestSuite.java | 2 +- .../server/mqtt/MqttSqlTestSuite.java | 2 +- .../server/rules/RuleEngineSqlTestSuite.java | 2 +- .../server/system/SystemSqlTestSuite.java | 2 +- .../SearchTextBasedWithAdditionalInfo.java | 3 +- .../common/data/kv/BaseAttributeKvEntry.java | 7 + .../server/common/data/kv/BasicKvEntry.java | 5 + .../server/common/data/kv/BasicTsKvEntry.java | 5 + .../server/common/data/kv/DataType.java | 2 +- .../server/common/data/kv/JsonDataEntry.java | 69 +++++ .../server/common/data/kv/KvEntry.java | 2 + .../transport/adaptor/JsonConverter.java | 40 ++- .../src/main/proto/transport.proto | 2 + .../CassandraBaseAttributesDao.java | 40 +-- .../server/dao/model/ModelConstants.java | 17 +- .../dao/model/sql/AbstractTsKvEntity.java | 8 +- .../dao/model/sql/AttributeKvEntity.java | 8 + .../dao/model/sqlts/hsql/TsKvEntity.java | 18 +- .../model/sqlts/latest/TsKvLatestEntity.java | 6 +- .../dao/model/sqlts/psql/TsKvEntity.java | 15 +- .../sqlts/timescale/TimescaleTsKvEntity.java | 17 +- .../AttributeKvInsertRepository.java | 118 ++------ .../HsqlAttributesInsertRepository.java | 34 +-- .../dao/sql/attributes/JpaAttributeDao.java | 1 + .../PsqlAttributesInsertRepository.java | 19 -- .../dao/sqlts/AbstractSqlTimeseriesDao.java | 2 + .../sqlts/hsql/HsqlInsertTsRepository.java | 12 +- .../dao/sqlts/hsql/TsKvHsqlRepository.java | 3 +- .../latest/HsqlLatestInsertTsRepository.java | 8 +- .../latest/PsqlLatestInsertTsRepository.java | 33 ++- .../latest/SearchTsKvLatestRepository.java | 3 +- .../sqlts/latest/TsKvLatestRepository.java | 4 - .../dao/sqlts/psql/JpaPsqlTimeseriesDao.java | 1 + .../sqlts/psql/PsqlInsertTsRepository.java | 21 +- .../dao/sqlts/psql/TsKvPsqlRepository.java | 3 +- .../timescale/AggregationRepository.java | 9 +- .../TimescaleInsertTsRepository.java | 19 +- .../timescale/TimescaleTimeseriesDao.java | 2 + .../AggregatePartitionsFunction.java | 48 +++- .../CassandraBaseTimeseriesDao.java | 59 ++-- .../resources/cassandra/schema-entities.cql | 1 + .../main/resources/cassandra/schema-ts.cql | 2 + .../resources/sql/schema-entities-hsql.sql | 251 ++++++++++++++++++ .../main/resources/sql/schema-entities.sql | 1 + .../main/resources/sql/schema-timescale.sql | 2 + dao/src/main/resources/sql/schema-ts-hsql.sql | 2 + dao/src/main/resources/sql/schema-ts-psql.sql | 4 +- .../server/dao/JpaDaoTestSuite.java | 2 +- .../server/dao/SqlDaoServiceTestSuite.java | 2 +- .../metadata/TbAbstractGetAttributesNode.java | 11 +- .../engine/metadata/TbGetTelemetryNode.java | 9 + 56 files changed, 713 insertions(+), 324 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/kv/JsonDataEntry.java create mode 100644 dao/src/main/resources/sql/schema-entities-hsql.sql 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 7eb91847cf..1be7e18aab 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 @@ -567,6 +567,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { case STRING_V: json.addProperty(kv.getKey(), kv.getStringV()); break; + case JSON_V: + json.add(kv.getKey(), jsonParser.parse(kv.getJsonV())); + break; } } return json; @@ -643,6 +646,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { builder.setType(KeyValueType.STRING_V); builder.setStringV(kvEntry.getStrValue().get()); break; + case JSON: + builder.setType(KeyValueType.JSON_V); + builder.setJsonV(kvEntry.getJsonValue().get()); + break; } return builder.build(); } diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 124a6e7ba1..634a78de27 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -642,6 +642,8 @@ public abstract class BaseController { entityNode.put(attr.getKey(), attr.getDoubleValue().get()); } else if (attr.getDataType() == DataType.LONG) { entityNode.put(attr.getKey(), attr.getLongValue().get()); + } else if (attr.getDataType() == DataType.JSON) { + entityNode.set(attr.getKey(), json.readTree(attr.getJsonValue().get())); } else { entityNode.put(attr.getKey(), attr.getValueAsString()); } diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index 43b525021e..86a2457e99 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java @@ -15,12 +15,15 @@ */ package org.thingsboard.server.controller; +import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.gson.JsonElement; +import com.google.gson.JsonParseException; import com.google.gson.JsonParser; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -56,8 +59,10 @@ import org.thingsboard.server.common.data.kv.BaseDeleteTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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.ReadTsKvQuery; @@ -77,6 +82,7 @@ import org.thingsboard.server.service.telemetry.exception.UncheckedApiException; import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; +import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; import java.util.HashSet; @@ -107,6 +113,8 @@ public class TelemetryController extends BaseController { private ExecutorService executor; + private static final ObjectMapper mapper = new ObjectMapper(); + @PostConstruct public void initExecutor() { executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("telemetry-controller")); @@ -284,8 +292,7 @@ public class TelemetryController extends BaseController { if (startTs == null || endTs == null) { deleteToTs = endTs; return getImmediateDeferredResult("When deleteAllDataForKeys is false, start and end timestamp values shouldn't be empty", HttpStatus.BAD_REQUEST); - } - else{ + } else { deleteFromTs = startTs; deleteToTs = endTs; } @@ -536,8 +543,9 @@ public class TelemetryController extends BaseController { return new FutureCallback>() { @Override public void onSuccess(List attributes) { - List values = attributes.stream().map(attribute -> new AttributeData(attribute.getLastUpdateTs(), - attribute.getKey(), attribute.getValue())).collect(Collectors.toList()); + List values = attributes.stream().map(attribute -> + new AttributeData(attribute.getLastUpdateTs(), attribute.getKey(), getKvValue(attribute)) + ).collect(Collectors.toList()); logAttributesRead(user, entityId, scope, keyList, null); response.setResult(new ResponseEntity<>(values, HttpStatus.OK)); } @@ -639,7 +647,9 @@ public class TelemetryController extends BaseController { jsonNode.fields().forEachRemaining(entry -> { String key = entry.getKey(); JsonNode value = entry.getValue(); - if (entry.getValue().isTextual()) { + if (entry.getValue().isObject() || entry.getValue().isArray()) { + attributes.add(new BaseAttributeKvEntry(new JsonDataEntry(key, toJsonStr(value)), ts)); + } else if (entry.getValue().isTextual()) { if (maxStringValueLength > 0 && entry.getValue().textValue().length() > maxStringValueLength) { String message = String.format("String value length [%d] for key [%s] is greater than maximum allowed [%d]", entry.getValue().textValue().length(), key, maxStringValueLength); throw new UncheckedApiException(new InvalidParametersException(message)); @@ -659,4 +669,27 @@ public class TelemetryController extends BaseController { }); return attributes; } + + private String toJsonStr(JsonNode value) { + try { + return mapper.writeValueAsString(value); + } catch (JsonProcessingException e) { + throw new JsonParseException("Can't parse jsonValue: " + value, e); + } + } + + private JsonNode toJsonNode(String value) { + try { + return mapper.readTree(value); + } catch (IOException e) { + throw new JsonParseException("Can't parse jsonValue: " + value, e); + } + } + + private Object getKvValue(KvEntry entry) { + if (entry.getDataType() == DataType.JSON) { + return toJsonNode(entry.getJsonValue().get()); + } + return entry.getValue(); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 45134a6bdd..d29e5de5b7 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -24,9 +24,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; +import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; -import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.server.actors.service.ActorService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; @@ -36,7 +36,20 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.*; +import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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.ReadTsKvQuery; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.cluster.SendToClusterMsg; import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.dao.attributes.AttributesService; @@ -105,7 +118,7 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio @Autowired @Lazy private ActorService actorService; - + private ExecutorService tsCallBackExecutor; private ExecutorService wsCallBackExecutor; @@ -692,6 +705,10 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio Optional doubleValue = attr.getDoubleValue(); doubleValue.ifPresent(dataBuilder::setDoubleValue); break; + case JSON: + Optional jsonValue = attr.getJsonValue(); + jsonValue.ifPresent(dataBuilder::setJsonValue); + break; case STRING: Optional stringValue = attr.getStrValue(); stringValue.ifPresent(dataBuilder::setStrValue); @@ -724,6 +741,9 @@ public class DefaultTelemetrySubscriptionService implements TelemetrySubscriptio case STRING: entry = new StringDataEntry(proto.getKey(), proto.getStrValue()); break; + case JSON: + entry = new JsonDataEntry(proto.getKey(), proto.getJsonValue()); + break; } return entry; } diff --git a/application/src/main/proto/cluster.proto b/application/src/main/proto/cluster.proto index b4ebc52f5e..cfacc66121 100644 --- a/application/src/main/proto/cluster.proto +++ b/application/src/main/proto/cluster.proto @@ -125,6 +125,7 @@ message KeyValueProto { int64 longValue = 5; double doubleValue = 6; bool boolValue = 7; + string jsonValue = 8; } message FromDeviceRPCResponseProto { diff --git a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java index 4fe33e4716..8dc0acff57 100644 --- a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java @@ -30,7 +30,7 @@ public class ControllerSqlTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/schema-entities-idx.sql", "sql/system-data.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/schema-entities-idx.sql", "sql/system-data.sql"), "sql/drop-all-tables.sql", "sql-test.properties"); } diff --git a/application/src/test/java/org/thingsboard/server/mqtt/MqttSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/mqtt/MqttSqlTestSuite.java index 5fb8c4d0c7..2863589ba1 100644 --- a/application/src/test/java/org/thingsboard/server/mqtt/MqttSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/mqtt/MqttSqlTestSuite.java @@ -29,7 +29,7 @@ public class MqttSqlTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/system-data.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/system-data.sql"), "sql/drop-all-tables.sql", "sql-test.properties"); } diff --git a/application/src/test/java/org/thingsboard/server/rules/RuleEngineSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/rules/RuleEngineSqlTestSuite.java index ce2c6852be..5f930821f7 100644 --- a/application/src/test/java/org/thingsboard/server/rules/RuleEngineSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/rules/RuleEngineSqlTestSuite.java @@ -30,7 +30,7 @@ public class RuleEngineSqlTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/system-data.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/system-data.sql"), "sql/drop-all-tables.sql", "sql-test.properties"); } diff --git a/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java index 3cbb7d9773..b12d513ce0 100644 --- a/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java @@ -31,7 +31,7 @@ public class SystemSqlTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/system-data.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/system-data.sql"), "sql/drop-all-tables.sql", "sql-test.properties"); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SearchTextBasedWithAdditionalInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/SearchTextBasedWithAdditionalInfo.java index ecbd2c73f5..8dc9bf6abc 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/SearchTextBasedWithAdditionalInfo.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SearchTextBasedWithAdditionalInfo.java @@ -35,6 +35,7 @@ import java.util.function.Consumer; @Slf4j public abstract class SearchTextBasedWithAdditionalInfo extends SearchTextBased implements HasAdditionalInfo { + private static final ObjectMapper mapper = new ObjectMapper(); private transient JsonNode additionalInfo; @JsonIgnore private byte[] additionalInfoBytes; @@ -97,7 +98,7 @@ public abstract class SearchTextBasedWithAdditionalInfo ext public static void setJson(JsonNode json, Consumer jsonConsumer, Consumer bytesConsumer) { jsonConsumer.accept(json); try { - bytesConsumer.accept(new ObjectMapper().writeValueAsBytes(json)); + bytesConsumer.accept(mapper.writeValueAsBytes(json)); } catch (JsonProcessingException e) { log.warn("Can't serialize json data: ", e); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseAttributeKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseAttributeKvEntry.java index ac8a5c2f78..5639f98d01 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseAttributeKvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseAttributeKvEntry.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.common.data.kv; +import com.fasterxml.jackson.databind.JsonNode; + import java.util.Optional; /** @@ -65,6 +67,11 @@ public class BaseAttributeKvEntry implements AttributeKvEntry { return kv.getDoubleValue(); } + @Override + public Optional getJsonValue() { + return kv.getJsonValue(); + } + @Override public String getValueAsString() { return kv.getValueAsString(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicKvEntry.java index a41bf6b232..7bc92ff74c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicKvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicKvEntry.java @@ -51,6 +51,11 @@ public abstract class BasicKvEntry implements KvEntry { return Optional.ofNullable(null); } + @Override + public Optional getJsonValue() { + return Optional.ofNullable(null); + } + @Override public boolean equals(Object o) { if (this == o) return true; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java index f7628da74c..c2d6688004 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java @@ -58,6 +58,11 @@ public class BasicTsKvEntry implements TsKvEntry { return kv.getDoubleValue(); } + @Override + public Optional getJsonValue() { + return kv.getJsonValue(); + } + @Override public Object getValue() { return kv.getValue(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/DataType.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/DataType.java index 84f918ede3..3571b0c882 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/DataType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/DataType.java @@ -17,6 +17,6 @@ package org.thingsboard.server.common.data.kv; public enum DataType { - STRING, LONG, BOOLEAN, DOUBLE; + STRING, LONG, BOOLEAN, DOUBLE, JSON; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/JsonDataEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/JsonDataEntry.java new file mode 100644 index 0000000000..0510f311d5 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/JsonDataEntry.java @@ -0,0 +1,69 @@ +/** + * 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.data.kv; + +import java.util.Objects; +import java.util.Optional; + +public class JsonDataEntry extends BasicKvEntry { + private final String value; + + public JsonDataEntry(String key, String value) { + super(key); + this.value = value; + } + + @Override + public DataType getDataType() { + return DataType.JSON; + } + + @Override + public Optional getJsonValue() { + return Optional.ofNullable(value); + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (!(o instanceof JsonDataEntry)) return false; + if (!super.equals(o)) return false; + JsonDataEntry that = (JsonDataEntry) o; + return Objects.equals(value, that.value); + } + + @Override + public Object getValue() { + return value; + } + + @Override + public int hashCode() { + return Objects.hash(super.hashCode(), value); + } + + @Override + public String toString() { + return "JsonDataEntry{" + + "value=" + value + + "} " + super.toString(); + } + + @Override + public String getValueAsString() { + return value; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java index c8753fda8b..296ddd37aa 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java @@ -37,6 +37,8 @@ public interface KvEntry extends Serializable { Optional getDoubleValue(); + Optional getJsonValue(); + String getValueAsString(); Object getValue(); 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 03f840a17e..56e2c87d6e 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 @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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; @@ -59,6 +60,7 @@ import java.util.stream.Collectors; public class JsonConverter { private static final Gson GSON = new Gson(); + private static final JsonParser JSON_PARSER = new JsonParser(); private static final String CAN_T_PARSE_VALUE = "Can't parse value: "; private static final String DEVICE_PROPERTY = "device"; @@ -204,6 +206,14 @@ public class JsonConverter { } else if (!value.isJsonNull()) { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + value); } + } else if (element.isJsonObject() || element.isJsonArray()) { + result.add(KeyValueProto + .newBuilder() + .setKey(valueEntry + .getKey()) + .setType(KeyValueType.JSON_V) + .setJsonV(element.toString()) + .build()); } else if (!element.isJsonNull()) { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + element); } @@ -354,6 +364,9 @@ public class JsonConverter { case LONG_V: json.addProperty(name, entry.getLongV()); break; + case JSON_V: + json.add(name, JSON_PARSER.parse(entry.getJsonV())); + break; } } @@ -363,47 +376,48 @@ public class JsonConverter { private static Consumer addToObjectFromProto(JsonObject result) { return de -> { - JsonPrimitive value; switch (de.getKv().getType()) { case BOOLEAN_V: - value = new JsonPrimitive(de.getKv().getBoolV()); + result.add(de.getKv().getKey(), new JsonPrimitive(de.getKv().getBoolV())); break; case DOUBLE_V: - value = new JsonPrimitive(de.getKv().getDoubleV()); + result.add(de.getKv().getKey(), new JsonPrimitive(de.getKv().getDoubleV())); break; case LONG_V: - value = new JsonPrimitive(de.getKv().getLongV()); + result.add(de.getKv().getKey(), new JsonPrimitive(de.getKv().getLongV())); break; case STRING_V: - value = new JsonPrimitive(de.getKv().getStringV()); + result.add(de.getKv().getKey(), new JsonPrimitive(de.getKv().getStringV())); break; + case JSON_V: + result.add(de.getKv().getKey(), JSON_PARSER.parse(de.getKv().getJsonV())); default: throw new IllegalArgumentException("Unsupported data type: " + de.getKv().getType()); } - result.add(de.getKv().getKey(), value); }; } private static Consumer addToObject(JsonObject result) { return de -> { - JsonPrimitive value; switch (de.getDataType()) { case BOOLEAN: - value = new JsonPrimitive(de.getBooleanValue().get()); + result.add(de.getKey(), new JsonPrimitive(de.getBooleanValue().get())); break; case DOUBLE: - value = new JsonPrimitive(de.getDoubleValue().get()); + result.add(de.getKey(), new JsonPrimitive(de.getDoubleValue().get())); break; case LONG: - value = new JsonPrimitive(de.getLongValue().get()); + result.add(de.getKey(), new JsonPrimitive(de.getLongValue().get())); break; case STRING: - value = new JsonPrimitive(de.getStrValue().get()); + result.add(de.getKey(), new JsonPrimitive(de.getStrValue().get())); + break; + case JSON: + result.add(de.getKey(), JSON_PARSER.parse(de.getJsonValue().get())); break; default: throw new IllegalArgumentException("Unsupported data type: " + de.getDataType()); } - result.add(de.getKey(), value); }; } @@ -464,6 +478,8 @@ public class JsonConverter { } else { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + value); } + } else if (element.isJsonObject() || element.isJsonArray()) { + result.add(new JsonDataEntry(valueEntry.getKey(), element.toString())); } else { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + element); } diff --git a/common/transport/transport-api/src/main/proto/transport.proto b/common/transport/transport-api/src/main/proto/transport.proto index 2d536b1769..e8b513574a 100644 --- a/common/transport/transport-api/src/main/proto/transport.proto +++ b/common/transport/transport-api/src/main/proto/transport.proto @@ -47,6 +47,7 @@ enum KeyValueType { LONG_V = 1; DOUBLE_V = 2; STRING_V = 3; + JSON_V = 4; } message KeyValueProto { @@ -56,6 +57,7 @@ message KeyValueProto { int64 long_v = 4; double double_v = 5; string string_v = 6; + string json_v = 7; } message TsKvProto { diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java index 93a39f8ef4..481f2ac085 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java @@ -112,31 +112,18 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) { - BoundStatement stmt = getSaveStmt().bind(); - stmt.setString(0, entityId.getEntityType().name()); - stmt.setUUID(1, entityId.getId()); - stmt.setString(2, attributeType); - stmt.setString(3, attribute.getKey()); - stmt.setLong(4, attribute.getLastUpdateTs()); - stmt.setString(5, attribute.getStrValue().orElse(null)); - Optional booleanValue = attribute.getBooleanValue(); - if (booleanValue.isPresent()) { - stmt.setBool(6, booleanValue.get()); - } else { - stmt.setToNull(6); - } - Optional longValue = attribute.getLongValue(); - if (longValue.isPresent()) { - stmt.setLong(7, longValue.get()); - } else { - stmt.setToNull(7); - } - Optional doubleValue = attribute.getDoubleValue(); - if (doubleValue.isPresent()) { - stmt.setDouble(8, doubleValue.get()); - } else { - stmt.setToNull(8); - } + BoundStatement stmt = getSaveStmt().bind() + .setString(0, entityId.getEntityType().name()) + .setUUID(1, entityId.getId()) + .setString(2, attributeType) + .setString(3, attribute.getKey()) + .setLong(4, attribute.getLastUpdateTs()) + .set(5, attribute.getStrValue().orElse(null), String.class) + .set(6, attribute.getBooleanValue().orElse(null), Boolean.class) + .set(7, attribute.getLongValue().orElse(null), Long.class) + .set(8, attribute.getDoubleValue().orElse(null), Double.class) + .set(9, attribute.getJsonValue().orElse(null), String.class); + log.trace("Generated save stmt [{}] for entityId {} and attributeType {} and attribute", stmt, entityId, attributeType, attribute); return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); } @@ -172,8 +159,9 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem "," + ModelConstants.BOOLEAN_VALUE_COLUMN + "," + ModelConstants.LONG_VALUE_COLUMN + "," + ModelConstants.DOUBLE_VALUE_COLUMN + + "," + ModelConstants.JSON_VALUE_COLUMN + ")" + - " VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?)"); + " VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"); } return saveStmt; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index a07c4868e6..96ce14c459 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -369,17 +369,18 @@ public class ModelConstants { public static final String STRING_VALUE_COLUMN = "str_v"; public static final String LONG_VALUE_COLUMN = "long_v"; public static final String DOUBLE_VALUE_COLUMN = "dbl_v"; + public static final String JSON_VALUE_COLUMN = "json_v"; - protected static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN}; + protected static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, JSON_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN}; - protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN)}; + protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN)}; - protected static final String[] MIN_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, - new String[]{min(LONG_VALUE_COLUMN), min(DOUBLE_VALUE_COLUMN), min(BOOLEAN_VALUE_COLUMN), min(STRING_VALUE_COLUMN)}); - protected static final String[] MAX_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, - new String[]{max(LONG_VALUE_COLUMN), max(DOUBLE_VALUE_COLUMN), max(BOOLEAN_VALUE_COLUMN), max(STRING_VALUE_COLUMN)}); - protected static final String[] SUM_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, - new String[]{sum(LONG_VALUE_COLUMN), sum(DOUBLE_VALUE_COLUMN)}); + protected static final String[] MIN_AGGREGATION_COLUMNS = + ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{min(LONG_VALUE_COLUMN), min(DOUBLE_VALUE_COLUMN), min(BOOLEAN_VALUE_COLUMN), min(STRING_VALUE_COLUMN), min(JSON_VALUE_COLUMN)}); + protected static final String[] MAX_AGGREGATION_COLUMNS = + ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{max(LONG_VALUE_COLUMN), max(DOUBLE_VALUE_COLUMN), max(BOOLEAN_VALUE_COLUMN), max(STRING_VALUE_COLUMN), max(JSON_VALUE_COLUMN)}); + protected static final String[] SUM_AGGREGATION_COLUMNS = + ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{sum(LONG_VALUE_COLUMN), sum(DOUBLE_VALUE_COLUMN)}); protected static final String[] AVG_AGGREGATION_COLUMNS = SUM_AGGREGATION_COLUMNS; public static String min(String s) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java index d7ffc72a34..f0dd03b5ca 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java @@ -19,6 +19,7 @@ import lombok.Data; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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; @@ -29,12 +30,12 @@ import javax.persistence.Column; import javax.persistence.Id; import javax.persistence.MappedSuperclass; import javax.persistence.Transient; - import java.util.UUID; import static org.thingsboard.server.dao.model.ModelConstants.BOOLEAN_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.DOUBLE_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.JSON_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.LONG_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.STRING_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.TS_COLUMN; @@ -68,6 +69,9 @@ public abstract class AbstractTsKvEntity implements ToData { @Column(name = DOUBLE_VALUE_COLUMN) protected Double doubleValue; + @Column(name = JSON_VALUE_COLUMN) + protected String jsonValue; + @Transient protected String strKey; @@ -93,6 +97,8 @@ public abstract class AbstractTsKvEntity implements ToData { kvEntry = new DoubleDataEntry(strKey, doubleValue); } else if (booleanValue != null) { kvEntry = new BooleanDataEntry(strKey, booleanValue); + } else if (jsonValue != null) { + kvEntry = new JsonDataEntry(strKey, jsonValue); } return new BasicTsKvEntry(ts, kvEntry); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AttributeKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AttributeKvEntity.java index 250d4325e5..f0de269cea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AttributeKvEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AttributeKvEntity.java @@ -20,6 +20,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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; @@ -33,6 +34,7 @@ import java.io.Serializable; import static org.thingsboard.server.dao.model.ModelConstants.BOOLEAN_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.DOUBLE_VALUE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.JSON_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.LAST_UPDATE_TS_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.LONG_VALUE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.STRING_VALUE_COLUMN; @@ -57,6 +59,9 @@ public class AttributeKvEntity implements ToData, Serializable @Column(name = DOUBLE_VALUE_COLUMN) private Double doubleValue; + @Column(name = JSON_VALUE_COLUMN) + private String jsonValue; + @Column(name = LAST_UPDATE_TS_COLUMN) private Long lastUpdateTs; @@ -71,7 +76,10 @@ public class AttributeKvEntity implements ToData, Serializable kvEntry = new DoubleDataEntry(id.getAttributeKey(), doubleValue); } else if (longValue != null) { kvEntry = new LongDataEntry(id.getAttributeKey(), longValue); + } else if (jsonValue != null) { + kvEntry = new JsonDataEntry(id.getAttributeKey(), jsonValue); } + return new BaseAttributeKvEntry(kvEntry, lastUpdateTs); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java index ba24543090..a38dc185a8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java @@ -16,30 +16,16 @@ package org.thingsboard.server.dao.model.sqlts.hsql; import lombok.Data; -import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.BooleanDataEntry; -import org.thingsboard.server.common.data.kv.DoubleDataEntry; -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.data.kv.TsKvEntry; import org.thingsboard.server.dao.model.ToData; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import javax.persistence.Column; import javax.persistence.Entity; -import javax.persistence.EnumType; -import javax.persistence.Enumerated; import javax.persistence.Id; import javax.persistence.IdClass; import javax.persistence.Table; -import javax.persistence.Transient; -import java.util.UUID; - -import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_TYPE_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN; @Data @@ -98,12 +84,14 @@ public final class TsKvEntity extends AbstractTsKvEntity implements ToData PATTERN_THREAD_LOCAL = ThreadLocal.withInitial(() -> Pattern.compile(String.valueOf(Character.MIN_VALUE))); private static final String EMPTY_STR = ""; - private static final String BATCH_UPDATE = "UPDATE attribute_kv SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, last_update_ts = ? " + + private static final String BATCH_UPDATE = "UPDATE attribute_kv SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ? " + "WHERE entity_type = ? and entity_id = ? and attribute_type =? and attribute_key = ?;"; private static final String INSERT_OR_UPDATE = - "INSERT INTO attribute_kv (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, last_update_ts) " + - "VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?) " + + "INSERT INTO attribute_kv (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, json_v, last_update_ts) " + + "VALUES(?, ?, ?, ?, ?, ?, ?, ?, cast(? AS json), ?) " + "ON CONFLICT (entity_type, entity_id, attribute_type, attribute_key) " + - "DO UPDATE SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, last_update_ts = ?;"; - - protected static final String BOOL_V = "bool_v"; - protected static final String STR_V = "str_v"; - protected static final String LONG_V = "long_v"; - protected static final String DBL_V = "dbl_v"; + "DO UPDATE SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ?;"; @Autowired protected JdbcTemplate jdbcTemplate; @@ -68,74 +60,6 @@ public abstract class AttributeKvInsertRepository { @Value("${sql.remove_null_chars}") private boolean removeNullChars; - @PersistenceContext - protected EntityManager entityManager; - - public abstract void saveOrUpdate(AttributeKvEntity entity); - - protected void processSaveOrUpdate(AttributeKvEntity entity, String requestBoolValue, String requestStrValue, String requestLongValue, String requestDblValue) { - if (entity.getBooleanValue() != null) { - saveOrUpdateBoolean(entity, requestBoolValue); - } - if (entity.getStrValue() != null) { - saveOrUpdateString(entity, requestStrValue); - } - if (entity.getLongValue() != null) { - saveOrUpdateLong(entity, requestLongValue); - } - if (entity.getDoubleValue() != null) { - saveOrUpdateDouble(entity, requestDblValue); - } - } - - @Modifying - private void saveOrUpdateBoolean(AttributeKvEntity entity, String query) { - entityManager.createNativeQuery(query) - .setParameter("entity_type", entity.getId().getEntityType().name()) - .setParameter("entity_id", entity.getId().getEntityId()) - .setParameter("attribute_type", entity.getId().getAttributeType()) - .setParameter("attribute_key", entity.getId().getAttributeKey()) - .setParameter("bool_v", entity.getBooleanValue()) - .setParameter("last_update_ts", entity.getLastUpdateTs()) - .executeUpdate(); - } - - @Modifying - private void saveOrUpdateString(AttributeKvEntity entity, String query) { - entityManager.createNativeQuery(query) - .setParameter("entity_type", entity.getId().getEntityType().name()) - .setParameter("entity_id", entity.getId().getEntityId()) - .setParameter("attribute_type", entity.getId().getAttributeType()) - .setParameter("attribute_key", entity.getId().getAttributeKey()) - .setParameter("str_v", replaceNullChars(entity.getStrValue())) - .setParameter("last_update_ts", entity.getLastUpdateTs()) - .executeUpdate(); - } - - @Modifying - private void saveOrUpdateLong(AttributeKvEntity entity, String query) { - entityManager.createNativeQuery(query) - .setParameter("entity_type", entity.getId().getEntityType().name()) - .setParameter("entity_id", entity.getId().getEntityId()) - .setParameter("attribute_type", entity.getId().getAttributeType()) - .setParameter("attribute_key", entity.getId().getAttributeKey()) - .setParameter("long_v", entity.getLongValue()) - .setParameter("last_update_ts", entity.getLastUpdateTs()) - .executeUpdate(); - } - - @Modifying - private void saveOrUpdateDouble(AttributeKvEntity entity, String query) { - entityManager.createNativeQuery(query) - .setParameter("entity_type", entity.getId().getEntityType().name()) - .setParameter("entity_id", entity.getId().getEntityId()) - .setParameter("attribute_type", entity.getId().getAttributeType()) - .setParameter("attribute_key", entity.getId().getAttributeKey()) - .setParameter("dbl_v", entity.getDoubleValue()) - .setParameter("last_update_ts", entity.getLastUpdateTs()) - .executeUpdate(); - } - protected void saveOrUpdate(List entities) { transactionTemplate.execute(new TransactionCallbackWithoutResult() { @Override @@ -164,11 +88,13 @@ public abstract class AttributeKvInsertRepository { ps.setNull(4, Types.BOOLEAN); } - ps.setLong(5, kvEntity.getLastUpdateTs()); - ps.setString(6, kvEntity.getId().getEntityType().name()); - ps.setString(7, kvEntity.getId().getEntityId()); - ps.setString(8, kvEntity.getId().getAttributeType()); - ps.setString(9, kvEntity.getId().getAttributeKey()); + ps.setString(5, replaceNullChars(kvEntity.getJsonValue())); + + ps.setLong(6, kvEntity.getLastUpdateTs()); + ps.setString(7, kvEntity.getId().getEntityType().name()); + ps.setString(8, kvEntity.getId().getEntityId()); + ps.setString(9, kvEntity.getId().getAttributeType()); + ps.setString(10, kvEntity.getId().getAttributeKey()); } @Override @@ -199,35 +125,39 @@ public abstract class AttributeKvInsertRepository { ps.setString(2, kvEntity.getId().getEntityId()); ps.setString(3, kvEntity.getId().getAttributeType()); ps.setString(4, kvEntity.getId().getAttributeKey()); + ps.setString(5, replaceNullChars(kvEntity.getStrValue())); - ps.setString(10, replaceNullChars(kvEntity.getStrValue())); + ps.setString(11, replaceNullChars(kvEntity.getStrValue())); if (kvEntity.getLongValue() != null) { ps.setLong(6, kvEntity.getLongValue()); - ps.setLong(11, kvEntity.getLongValue()); + ps.setLong(12, kvEntity.getLongValue()); } else { ps.setNull(6, Types.BIGINT); - ps.setNull(11, Types.BIGINT); + ps.setNull(12, Types.BIGINT); } if (kvEntity.getDoubleValue() != null) { ps.setDouble(7, kvEntity.getDoubleValue()); - ps.setDouble(12, kvEntity.getDoubleValue()); + ps.setDouble(13, kvEntity.getDoubleValue()); } else { ps.setNull(7, Types.DOUBLE); - ps.setNull(12, Types.DOUBLE); + ps.setNull(13, Types.DOUBLE); } if (kvEntity.getBooleanValue() != null) { ps.setBoolean(8, kvEntity.getBooleanValue()); - ps.setBoolean(13, kvEntity.getBooleanValue()); + ps.setBoolean(14, kvEntity.getBooleanValue()); } else { ps.setNull(8, Types.BOOLEAN); - ps.setNull(13, Types.BOOLEAN); + ps.setNull(14, Types.BOOLEAN); } - ps.setLong(9, kvEntity.getLastUpdateTs()); - ps.setLong(14, kvEntity.getLastUpdateTs()); + ps.setString(9, replaceNullChars(kvEntity.getJsonValue())); + ps.setString(15, replaceNullChars(kvEntity.getJsonValue())); + + ps.setLong(10, kvEntity.getLastUpdateTs()); + ps.setLong(16, kvEntity.getLastUpdateTs()); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/HsqlAttributesInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/HsqlAttributesInsertRepository.java index 8378a4488b..1d7aad384a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/HsqlAttributesInsertRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/HsqlAttributesInsertRepository.java @@ -30,36 +30,16 @@ import java.util.List; @Transactional public class HsqlAttributesInsertRepository extends AttributeKvInsertRepository { - private static final String ON_BOOL_VALUE_UPDATE_SET_NULLS = " attribute_kv.str_v = null, attribute_kv.long_v = null, attribute_kv.dbl_v = null "; - private static final String ON_STR_VALUE_UPDATE_SET_NULLS = " attribute_kv.bool_v = null, attribute_kv.long_v = null, attribute_kv.dbl_v = null "; - private static final String ON_LONG_VALUE_UPDATE_SET_NULLS = " attribute_kv.str_v = null, attribute_kv.bool_v = null, attribute_kv.dbl_v = null "; - private static final String ON_DBL_VALUE_UPDATE_SET_NULLS = " attribute_kv.str_v = null, attribute_kv.long_v = null, attribute_kv.bool_v = null "; - - private static final String INSERT_BOOL_STATEMENT = getInsertOrUpdateString(BOOL_V, ON_BOOL_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_STR_STATEMENT = getInsertOrUpdateString(STR_V, ON_STR_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_LONG_STATEMENT = getInsertOrUpdateString(LONG_V, ON_LONG_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_DBL_STATEMENT = getInsertOrUpdateString(DBL_V, ON_DBL_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_OR_UPDATE = - "MERGE INTO attribute_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?, ?) " + - "A (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, last_update_ts) " + + "MERGE INTO attribute_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) " + + "A (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, json_v, last_update_ts) " + "ON (attribute_kv.entity_type=A.entity_type " + "AND attribute_kv.entity_id=A.entity_id " + "AND attribute_kv.attribute_type=A.attribute_type " + "AND attribute_kv.attribute_key=A.attribute_key) " + - "WHEN MATCHED THEN UPDATE SET attribute_kv.str_v = A.str_v, attribute_kv.long_v = A.long_v, attribute_kv.dbl_v = A.dbl_v, attribute_kv.bool_v = A.bool_v, attribute_kv.last_update_ts = A.last_update_ts " + - "WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, last_update_ts) " + - "VALUES (A.entity_type, A.entity_id, A.attribute_type, A.attribute_key, A.str_v, A.long_v, A.dbl_v, A.bool_v, A.last_update_ts)"; - - @Override - public void saveOrUpdate(AttributeKvEntity entity) { - processSaveOrUpdate(entity, INSERT_BOOL_STATEMENT, INSERT_STR_STATEMENT, INSERT_LONG_STATEMENT, INSERT_DBL_STATEMENT); - } - - private static String getInsertOrUpdateString(String value, String nullValues) { - return "MERGE INTO attribute_kv USING(VALUES :entity_type, :entity_id, :attribute_type, :attribute_key, :" + value + ", :last_update_ts) A (entity_type, entity_id, attribute_type, attribute_key, " + value + ", last_update_ts) ON (attribute_kv.entity_type=A.entity_type AND attribute_kv.entity_id=A.entity_id AND attribute_kv.attribute_type=A.attribute_type AND attribute_kv.attribute_key=A.attribute_key) WHEN MATCHED THEN UPDATE SET attribute_kv." + value + " = A." + value + ", attribute_kv.last_update_ts = A.last_update_ts," + nullValues + "WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, attribute_type, attribute_key, " + value + ", last_update_ts) VALUES (A.entity_type, A.entity_id, A.attribute_type, A.attribute_key, A." + value + ", A.last_update_ts)"; - } - + "WHEN MATCHED THEN UPDATE SET attribute_kv.str_v = A.str_v, attribute_kv.long_v = A.long_v, attribute_kv.dbl_v = A.dbl_v, attribute_kv.bool_v = A.bool_v, attribute_kv.json_v = A.json_v, attribute_kv.last_update_ts = A.last_update_ts " + + "WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, json_v, last_update_ts) " + + "VALUES (A.entity_type, A.entity_id, A.attribute_type, A.attribute_key, A.str_v, A.long_v, A.dbl_v, A.bool_v, A.json_v, A.last_update_ts)"; @Override protected void saveOrUpdate(List entities) { @@ -89,7 +69,9 @@ public class HsqlAttributesInsertRepository extends AttributeKvInsertRepository ps.setNull(8, Types.BOOLEAN); } - ps.setLong(9, entity.getLastUpdateTs()); + ps.setString(9, entity.getJsonValue()); + + ps.setLong(10, entity.getLastUpdateTs()); }); }); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index 85a4541be2..ab9964d3e3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java @@ -128,6 +128,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl entity.setDoubleValue(attribute.getDoubleValue().orElse(null)); entity.setLongValue(attribute.getLongValue().orElse(null)); entity.setBooleanValue(attribute.getBooleanValue().orElse(null)); + entity.setJsonValue(attribute.getJsonValue().orElse(null)); return addToQueue(entity); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/PsqlAttributesInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/PsqlAttributesInsertRepository.java index 1f553ada04..020e63cd36 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/PsqlAttributesInsertRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/PsqlAttributesInsertRepository.java @@ -17,7 +17,6 @@ package org.thingsboard.server.dao.sql.attributes; import org.springframework.stereotype.Repository; import org.springframework.transaction.annotation.Transactional; -import org.thingsboard.server.dao.model.sql.AttributeKvEntity; import org.thingsboard.server.dao.util.PsqlDao; import org.thingsboard.server.dao.util.SqlDao; @@ -27,22 +26,4 @@ import org.thingsboard.server.dao.util.SqlDao; @Transactional public class PsqlAttributesInsertRepository extends AttributeKvInsertRepository { - private static final String ON_BOOL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, dbl_v = null"; - private static final String ON_STR_VALUE_UPDATE_SET_NULLS = "bool_v = null, long_v = null, dbl_v = null"; - private static final String ON_LONG_VALUE_UPDATE_SET_NULLS = "str_v = null, bool_v = null, dbl_v = null"; - private static final String ON_DBL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, bool_v = null"; - - private static final String INSERT_OR_UPDATE_BOOL_STATEMENT = getInsertOrUpdateString(BOOL_V, ON_BOOL_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_OR_UPDATE_STR_STATEMENT = getInsertOrUpdateString(STR_V, ON_STR_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_OR_UPDATE_LONG_STATEMENT = getInsertOrUpdateString(LONG_V , ON_LONG_VALUE_UPDATE_SET_NULLS); - private static final String INSERT_OR_UPDATE_DBL_STATEMENT = getInsertOrUpdateString(DBL_V, ON_DBL_VALUE_UPDATE_SET_NULLS); - - @Override - public void saveOrUpdate(AttributeKvEntity entity) { - processSaveOrUpdate(entity, INSERT_OR_UPDATE_BOOL_STATEMENT, INSERT_OR_UPDATE_STR_STATEMENT, INSERT_OR_UPDATE_LONG_STATEMENT, INSERT_OR_UPDATE_DBL_STATEMENT); - } - - private static String getInsertOrUpdateString(String value, String nullValues) { - return "INSERT INTO attribute_kv (entity_type, entity_id, attribute_type, attribute_key, " + value + ", last_update_ts) VALUES (:entity_type, :entity_id, :attribute_type, :attribute_key, :" + value + ", :last_update_ts) ON CONFLICT (entity_type, entity_id, attribute_type, attribute_key) DO UPDATE SET " + value + " = :" + value + ", last_update_ts = :last_update_ts," + nullValues; - } } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java index 63893c3a5e..e8cbf5d73b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java @@ -253,6 +253,8 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx latestEntity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null)); latestEntity.setLongValue(tsKvEntry.getLongValue().orElse(null)); latestEntity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null)); + latestEntity.setJsonValue(tsKvEntry.getJsonValue().orElse(null)); + return tsLatestQueue.add(latestEntity); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java index 2d10f35cd9..d1e5294309 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java @@ -37,14 +37,14 @@ import java.util.List; public class HsqlInsertTsRepository extends AbstractInsertRepository implements InsertTsRepository { private static final String INSERT_OR_UPDATE = - "MERGE INTO ts_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?) " + - "T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " + + "MERGE INTO ts_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?) " + + "T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) " + "ON (ts_kv.entity_id=T.entity_id " + "AND ts_kv.key=T.key " + "AND ts_kv.ts=T.ts) " + - "WHEN MATCHED THEN UPDATE SET ts_kv.bool_v = T.bool_v, ts_kv.str_v = T.str_v, ts_kv.long_v = T.long_v, ts_kv.dbl_v = T.dbl_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 MATCHED THEN UPDATE SET ts_kv.bool_v = T.bool_v, ts_kv.str_v = T.str_v, ts_kv.long_v = T.long_v, ts_kv.dbl_v = T.dbl_v ,ts_kv.json_v = T.json_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) { @@ -76,6 +76,8 @@ public class HsqlInsertTsRepository extends AbstractInsertRepository implements } else { ps.setNull(7, Types.DOUBLE); } + + ps.setString(8, tsKvEntity.getJsonValue()); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java index 552d515e44..a7c0effb97 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java @@ -97,7 +97,8 @@ public interface TsKvHsqlRepository extends CrudRepository :startTs AND tskv.ts <= :endTs") CompletableFuture findCount(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java index 931eb86689..65ac6257f0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java @@ -36,11 +36,11 @@ import java.util.List; public class HsqlLatestInsertTsRepository extends AbstractInsertRepository implements InsertLatestTsRepository { private static final String INSERT_OR_UPDATE = - "MERGE INTO ts_kv_latest USING(VALUES ?, ?, ?, ?, ?, ?, ?) " + - "T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " + + "MERGE INTO ts_kv_latest USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?) " + + "T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) " + "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 " + + "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);"; @@ -72,6 +72,8 @@ public class HsqlLatestInsertTsRepository extends AbstractInsertRepository imple } else { ps.setNull(7, Types.DOUBLE); } + + ps.setString(8, entities.get(i).getJsonValue()); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java index 16a4c4ccd8..d367f44620 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java @@ -38,12 +38,12 @@ import java.util.List; public class PsqlLatestInsertTsRepository extends AbstractInsertRepository implements InsertLatestTsRepository { private static final String BATCH_UPDATE = - "UPDATE ts_kv_latest SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ? WHERE entity_id = ? and key = ?"; + "UPDATE ts_kv_latest SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?, json_v = cast(? AS json) WHERE entity_id = ? and key = ?"; private static final String INSERT_OR_UPDATE = - "INSERT INTO ts_kv_latest (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES(?, ?, ?, ?, ?, ?, ?) " + - "ON CONFLICT (entity_id, key) DO UPDATE SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?;"; + "INSERT INTO ts_kv_latest (entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) VALUES(?, ?, ?, ?, ?, ?, ?, cast(? AS json)) " + + "ON CONFLICT (entity_id, key) DO UPDATE SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?, json_v = cast(? AS json);"; @Override public void saveOrUpdate(List entities) { @@ -76,8 +76,10 @@ public class PsqlLatestInsertTsRepository extends AbstractInsertRepository imple ps.setNull(5, Types.DOUBLE); } - ps.setObject(6, tsKvLatestEntity.getEntityId()); - ps.setInt(7, tsKvLatestEntity.getKey()); + ps.setString(6, replaceNullChars(tsKvLatestEntity.getJsonValue())); + + ps.setObject(7, tsKvLatestEntity.getEntityId()); + ps.setInt(8, tsKvLatestEntity.getKey()); } @Override @@ -106,36 +108,39 @@ public class PsqlLatestInsertTsRepository extends AbstractInsertRepository imple TsKvLatestEntity tsKvLatestEntity = insertEntities.get(i); ps.setObject(1, tsKvLatestEntity.getEntityId()); ps.setInt(2, tsKvLatestEntity.getKey()); + ps.setLong(3, tsKvLatestEntity.getTs()); - ps.setLong(8, tsKvLatestEntity.getTs()); + ps.setLong(9, tsKvLatestEntity.getTs()); if (tsKvLatestEntity.getBooleanValue() != null) { ps.setBoolean(4, tsKvLatestEntity.getBooleanValue()); - ps.setBoolean(9, tsKvLatestEntity.getBooleanValue()); + ps.setBoolean(10, tsKvLatestEntity.getBooleanValue()); } else { ps.setNull(4, Types.BOOLEAN); - ps.setNull(9, Types.BOOLEAN); + ps.setNull(10, Types.BOOLEAN); } ps.setString(5, replaceNullChars(tsKvLatestEntity.getStrValue())); - ps.setString(10, replaceNullChars(tsKvLatestEntity.getStrValue())); - + ps.setString(11, replaceNullChars(tsKvLatestEntity.getStrValue())); if (tsKvLatestEntity.getLongValue() != null) { ps.setLong(6, tsKvLatestEntity.getLongValue()); - ps.setLong(11, tsKvLatestEntity.getLongValue()); + ps.setLong(12, tsKvLatestEntity.getLongValue()); } else { ps.setNull(6, Types.BIGINT); - ps.setNull(11, Types.BIGINT); + ps.setNull(12, Types.BIGINT); } if (tsKvLatestEntity.getDoubleValue() != null) { ps.setDouble(7, tsKvLatestEntity.getDoubleValue()); - ps.setDouble(12, tsKvLatestEntity.getDoubleValue()); + ps.setDouble(13, tsKvLatestEntity.getDoubleValue()); } else { ps.setNull(7, Types.DOUBLE); - ps.setNull(12, Types.DOUBLE); + ps.setNull(13, Types.DOUBLE); } + + ps.setString(8, replaceNullChars(tsKvLatestEntity.getJsonValue())); + ps.setString(14, replaceNullChars(tsKvLatestEntity.getJsonValue())); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java index 5940a33d31..b5da093b6f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java @@ -29,8 +29,9 @@ import java.util.UUID; public class SearchTsKvLatestRepository { public static final String FIND_ALL_BY_ENTITY_ID = "findAllByEntityId"; + public static final String FIND_ALL_BY_ENTITY_ID_QUERY = "SELECT ts_kv_latest.entity_id AS entityId, ts_kv_latest.key AS key, ts_kv_dictionary.key AS strKey, ts_kv_latest.str_v AS strValue," + - " ts_kv_latest.bool_v AS boolValue, ts_kv_latest.long_v AS longValue, ts_kv_latest.dbl_v AS doubleValue, ts_kv_latest.ts AS ts FROM ts_kv_latest " + + " ts_kv_latest.bool_v AS boolValue, ts_kv_latest.long_v AS longValue, ts_kv_latest.dbl_v AS doubleValue, ts_kv_latest.json_v AS jsonValue, ts_kv_latest.ts AS ts FROM ts_kv_latest " + "INNER JOIN ts_kv_dictionary ON ts_kv_latest.key = ts_kv_dictionary.key_id WHERE ts_kv_latest.entity_id = cast(:id AS uuid)"; @PersistenceContext diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java index 9ba59c10ef..5232b88489 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java @@ -20,11 +20,7 @@ import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestCompositeKey; import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity; import org.thingsboard.server.dao.util.SqlDao; -import java.util.List; -import java.util.UUID; - @SqlDao public interface TsKvLatestRepository extends CrudRepository { - List findAllByEntityId(UUID entityId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java index f598742713..6d60e843ec 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java @@ -104,6 +104,7 @@ public class JpaPsqlTimeseriesDao 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)); PsqlPartition psqlPartition = toPartition(tsKvEntry.getTs()); log.trace("Saving entity: {}", entity); return tsQueue.add(new EntityContainer(entity, psqlPartition.getPartitionDate())); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java index e9cf5c9b03..caf4528812 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java @@ -41,8 +41,8 @@ public class PsqlInsertTsRepository extends AbstractInsertRepository implements private static final String INSERT_INTO_TS_KV = "INSERT INTO ts_kv_"; - private static final String VALUES_ON_CONFLICT_DO_UPDATE = " (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES (?, ?, ?, ?, ?, ?, ?) " + - "ON CONFLICT (entity_id, key, ts) DO UPDATE SET bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?;"; + private static final String VALUES_ON_CONFLICT_DO_UPDATE = " (entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) VALUES (?, ?, ?, ?, ?, ?, ?, cast(? AS json)) " + + "ON CONFLICT (entity_id, key, ts) DO UPDATE SET bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?, json_v = cast(? AS json);"; @Override public void saveOrUpdate(List> entities) { @@ -61,30 +61,33 @@ public class PsqlInsertTsRepository extends AbstractInsertRepository implements if (tsKvEntity.getBooleanValue() != null) { ps.setBoolean(4, tsKvEntity.getBooleanValue()); - ps.setBoolean(8, tsKvEntity.getBooleanValue()); + ps.setBoolean(9, tsKvEntity.getBooleanValue()); } else { ps.setNull(4, Types.BOOLEAN); - ps.setNull(8, Types.BOOLEAN); + ps.setNull(9, Types.BOOLEAN); } ps.setString(5, replaceNullChars(tsKvEntity.getStrValue())); - ps.setString(9, replaceNullChars(tsKvEntity.getStrValue())); + ps.setString(10, replaceNullChars(tsKvEntity.getStrValue())); if (tsKvEntity.getLongValue() != null) { ps.setLong(6, tsKvEntity.getLongValue()); - ps.setLong(10, tsKvEntity.getLongValue()); + ps.setLong(11, tsKvEntity.getLongValue()); } else { ps.setNull(6, Types.BIGINT); - ps.setNull(10, Types.BIGINT); + ps.setNull(11, Types.BIGINT); } if (tsKvEntity.getDoubleValue() != null) { ps.setDouble(7, tsKvEntity.getDoubleValue()); - ps.setDouble(11, tsKvEntity.getDoubleValue()); + ps.setDouble(12, tsKvEntity.getDoubleValue()); } else { ps.setNull(7, Types.DOUBLE); - ps.setNull(11, Types.DOUBLE); + ps.setNull(12, Types.DOUBLE); + + ps.setString(8, replaceNullChars(tsKvEntity.getJsonValue())); + ps.setString(13, replaceNullChars(tsKvEntity.getJsonValue())); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java index 7b3328f86b..b164046bbc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java @@ -98,7 +98,8 @@ public interface TsKvPsqlRepository extends CrudRepository :startTs AND tskv.ts <= :endTs") CompletableFuture findCount(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java index 15deb702a7..5a0b9c6a59 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java @@ -36,18 +36,17 @@ public class AggregationRepository { public static final String FIND_SUM = "findSum"; public static final String FIND_COUNT = "findCount"; - public static final String FROM_WHERE_CLAUSE = "FROM tenant_ts_kv tskv WHERE tskv.tenant_id = cast(:tenantId AS uuid) AND tskv.entity_id = cast(:entityId AS uuid) AND tskv.key= cast(:entityKey AS int) AND tskv.ts > :startTs AND tskv.ts <= :endTs GROUP BY tskv.tenant_id, tskv.entity_id, tskv.key, tsBucket ORDER BY tskv.tenant_id, tskv.entity_id, tskv.key, tsBucket"; public static final String FIND_AVG_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, 'AVG' AS aggType "; - public static final String FIND_MAX_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MAX(COALESCE(tskv.long_v, -9223372036854775807)) AS longValue, MAX(COALESCE(tskv.dbl_v, -1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MAX(tskv.str_v) AS strValue, 'MAX' AS aggType "; + public static final String FIND_MAX_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MAX(COALESCE(tskv.long_v, -9223372036854775807)) AS longValue, MAX(COALESCE(tskv.dbl_v, -1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MAX(tskv.str_v) AS strValue, MAX(tskv.json_v) AS jsonValue, 'MAX' AS aggType "; - public static final String FIND_MIN_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MIN(COALESCE(tskv.long_v, 9223372036854775807)) AS longValue, MIN(COALESCE(tskv.dbl_v, 1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MIN(tskv.str_v) AS strValue, 'MIN' AS aggType "; + public static final String FIND_MIN_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MIN(COALESCE(tskv.long_v, 9223372036854775807)) AS longValue, MIN(COALESCE(tskv.dbl_v, 1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MIN(tskv.str_v) AS strValue, MIN(tskv.json_v) AS jsonValue,'MIN' AS aggType "; - public static final String FIND_SUM_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, 'SUM' AS aggType "; + public static final String FIND_SUM_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, null AS jsonValue, 'SUM' AS aggType "; - public static final String FIND_COUNT_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(CASE WHEN tskv.bool_v IS NULL THEN 0 ELSE 1 END) AS booleanValueCount, SUM(CASE WHEN tskv.str_v IS NULL THEN 0 ELSE 1 END) AS strValueCount, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longValueCount, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleValueCount "; + public static final String FIND_COUNT_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(CASE WHEN tskv.bool_v IS NULL THEN 0 ELSE 1 END) AS booleanValueCount, SUM(CASE WHEN tskv.str_v IS NULL THEN 0 ELSE 1 END) AS strValueCount, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longValueCount, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleValueCount, SUM(CASE WHEN tskv.json_v IS NULL THEN 0 ELSE 1 END) AS jsonValueCount "; @PersistenceContext private EntityManager entityManager; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java index a6fad65fbf..6d863af105 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java @@ -37,8 +37,8 @@ import java.util.List; public class TimescaleInsertTsRepository extends AbstractInsertRepository implements InsertTsRepository { private static final String INSERT_OR_UPDATE = - "INSERT INTO tenant_ts_kv (tenant_id, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES(?, ?, ?, ?, ?, ?, ?, ?) " + - "ON CONFLICT (tenant_id, entity_id, key, ts) DO UPDATE SET bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?;"; + "INSERT INTO tenant_ts_kv (tenant_id, entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) VALUES(?, ?, ?, ?, ?, ?, ?, ?, cast(? AS json)) " + + "ON CONFLICT (tenant_id, entity_id, key, ts) DO UPDATE SET bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?, json_v = cast(? AS json);"; @Override public void saveOrUpdate(List> entities) { @@ -56,28 +56,31 @@ public class TimescaleInsertTsRepository extends AbstractInsertRepository implem ps.setBoolean(9, tsKvEntity.getBooleanValue()); } else { ps.setNull(5, Types.BOOLEAN); - ps.setNull(9, Types.BOOLEAN); + ps.setNull(10, Types.BOOLEAN); } ps.setString(6, replaceNullChars(tsKvEntity.getStrValue())); - ps.setString(10, replaceNullChars(tsKvEntity.getStrValue())); + ps.setString(11, replaceNullChars(tsKvEntity.getStrValue())); if (tsKvEntity.getLongValue() != null) { ps.setLong(7, tsKvEntity.getLongValue()); - ps.setLong(11, tsKvEntity.getLongValue()); + ps.setLong(12, tsKvEntity.getLongValue()); } else { ps.setNull(7, Types.BIGINT); - ps.setNull(11, Types.BIGINT); + ps.setNull(12, Types.BIGINT); } if (tsKvEntity.getDoubleValue() != null) { ps.setDouble(8, tsKvEntity.getDoubleValue()); - ps.setDouble(12, tsKvEntity.getDoubleValue()); + ps.setDouble(13, tsKvEntity.getDoubleValue()); } else { ps.setNull(8, Types.DOUBLE); - ps.setNull(12, Types.DOUBLE); + ps.setNull(13, Types.DOUBLE); } + + ps.setString(9, replaceNullChars(tsKvEntity.getJsonValue())); + ps.setString(14, replaceNullChars(tsKvEntity.getJsonValue())); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java index f7e51ef18c..9f8f5c6f74 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java @@ -174,6 +174,8 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements 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 to timescale db: {}", entity); return tsQueue.add(new EntityContainer(entity, null)); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java index d81ddc3801..6229f68818 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -41,10 +42,12 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct private static final int DOUBLE_CNT_POS = 1; private static final int BOOL_CNT_POS = 2; private static final int STR_CNT_POS = 3; - private static final int LONG_POS = 4; - private static final int DOUBLE_POS = 5; - private static final int BOOL_POS = 6; - private static final int STR_POS = 7; + private static final int JSON_CNT_POS = 4; + private static final int LONG_POS = 5; + private static final int DOUBLE_POS = 6; + private static final int BOOL_POS = 7; + private static final int STR_POS = 8; + private static final int JSON_POS = 9; private final Aggregation aggregation; private final String key; @@ -72,7 +75,7 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct } } return processAggregationResult(aggResult); - }catch (Exception e){ + } catch (Exception e) { log.error("[{}][{}][{}] Failed to aggregate data", key, ts, aggregation, e); return Optional.empty(); } @@ -85,11 +88,13 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct Double curDValue = null; Boolean curBValue = null; String curSValue = null; + String curJValue = null; long longCount = row.getLong(LONG_CNT_POS); long doubleCount = row.getLong(DOUBLE_CNT_POS); long boolCount = row.getLong(BOOL_CNT_POS); long strCount = row.getLong(STR_CNT_POS); + long jsonCount = row.getLong(JSON_CNT_POS); if (longCount > 0 || doubleCount > 0) { if (longCount > 0) { @@ -111,6 +116,10 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct aggResult.dataType = DataType.STRING; curCount = strCount; curSValue = getStringValue(row); + } else if (jsonCount > 0) { + aggResult.dataType = DataType.JSON; + curCount = jsonCount; + curJValue = getJsonValue(row); } else { return; } @@ -120,9 +129,9 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { processAvgOrSumAggregation(aggResult, curCount, curLValue, curDValue); } else if (aggregation == Aggregation.MIN) { - processMinAggregation(aggResult, curLValue, curDValue, curBValue, curSValue); + processMinAggregation(aggResult, curLValue, curDValue, curBValue, curSValue, curJValue); } else if (aggregation == Aggregation.MAX) { - processMaxAggregation(aggResult, curLValue, curDValue, curBValue, curSValue); + processMaxAggregation(aggResult, curLValue, curDValue, curBValue, curSValue, curJValue); } } @@ -136,7 +145,7 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct } } - private void processMinAggregation(AggregationResult aggResult, Long curLValue, Double curDValue, Boolean curBValue, String curSValue) { + private void processMinAggregation(AggregationResult aggResult, Long curLValue, Double curDValue, Boolean curBValue, String curSValue, String curJValue) { if (curDValue != null || curLValue != null) { if (curDValue != null) { aggResult.dValue = aggResult.dValue == null ? curDValue : Math.min(aggResult.dValue, curDValue); @@ -148,10 +157,12 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct aggResult.bValue = aggResult.bValue == null ? curBValue : aggResult.bValue && curBValue; } else if (curSValue != null && (aggResult.sValue == null || curSValue.compareTo(aggResult.sValue) < 0)) { aggResult.sValue = curSValue; + } else if (curJValue != null && (aggResult.jValue == null || curJValue.compareTo(aggResult.jValue) < 0)) { + aggResult.jValue = curJValue; } } - private void processMaxAggregation(AggregationResult aggResult, Long curLValue, Double curDValue, Boolean curBValue, String curSValue) { + private void processMaxAggregation(AggregationResult aggResult, Long curLValue, Double curDValue, Boolean curBValue, String curSValue, String curJValue) { if (curDValue != null || curLValue != null) { if (curDValue != null) { aggResult.dValue = aggResult.dValue == null ? curDValue : Math.max(aggResult.dValue, curDValue); @@ -163,6 +174,8 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct aggResult.bValue = aggResult.bValue == null ? curBValue : aggResult.bValue || curBValue; } else if (curSValue != null && (aggResult.sValue == null || curSValue.compareTo(aggResult.sValue) > 0)) { aggResult.sValue = curSValue; + } else if (curJValue != null && (aggResult.jValue == null || curJValue.compareTo(aggResult.jValue) > 0)) { + aggResult.jValue = curJValue; } } @@ -182,6 +195,14 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct } } + private String getJsonValue(Row row) { + if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { + return row.getString(JSON_POS); + } else { + return null; + } + } + private Long getLongValue(Row row) { if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX || aggregation == Aggregation.SUM || aggregation == Aggregation.AVG) { @@ -223,7 +244,7 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct if (aggResult.count == 0 || (aggResult.dataType == DataType.DOUBLE && aggResult.dValue == null) || (aggResult.dataType == DataType.LONG && aggResult.lValue == null)) { return Optional.empty(); } else if (aggResult.dataType == DataType.DOUBLE || aggResult.dataType == DataType.LONG) { - if(aggregation == Aggregation.AVG || aggResult.hasDouble) { + if (aggregation == Aggregation.AVG || aggResult.hasDouble) { double sum = Optional.ofNullable(aggResult.dValue).orElse(0.0d) + Optional.ofNullable(aggResult.lValue).orElse(0L); return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, aggregation == Aggregation.SUM ? sum : (sum / aggResult.count)))); } else { @@ -235,15 +256,17 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct private Optional processMinOrMaxResult(AggregationResult aggResult) { if (aggResult.dataType == DataType.DOUBLE || aggResult.dataType == DataType.LONG) { - if(aggResult.hasDouble) { + if (aggResult.hasDouble) { double currentD = aggregation == Aggregation.MIN ? Optional.ofNullable(aggResult.dValue).orElse(Double.MAX_VALUE) : Optional.ofNullable(aggResult.dValue).orElse(Double.MIN_VALUE); double currentL = aggregation == Aggregation.MIN ? Optional.ofNullable(aggResult.lValue).orElse(Long.MAX_VALUE) : Optional.ofNullable(aggResult.lValue).orElse(Long.MIN_VALUE); return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, aggregation == Aggregation.MIN ? Math.min(currentD, currentL) : Math.max(currentD, currentL)))); } else { return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggResult.lValue))); } - } else if (aggResult.dataType == DataType.STRING) { + } else if (aggResult.dataType == DataType.STRING) { return Optional.of(new BasicTsKvEntry(ts, new StringDataEntry(key, aggResult.sValue))); + } else if (aggResult.dataType == DataType.JSON) { + return Optional.of(new BasicTsKvEntry(ts, new JsonDataEntry(key, aggResult.jValue))); } else { return Optional.of(new BasicTsKvEntry(ts, new BooleanDataEntry(key, aggResult.bValue))); } @@ -253,6 +276,7 @@ public class AggregatePartitionsFunction implements com.google.common.base.Funct DataType dataType = null; Boolean bValue = null; String sValue = null; + String jValue = null; Double dValue = null; Long lValue = null; long count = 0; diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index d0ebb49572..8fc8b4ab8a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -29,6 +29,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.core.env.Environment; @@ -42,6 +43,7 @@ import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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.ReadTsKvQuery; @@ -337,21 +339,31 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.BOOLEAN)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.DOUBLE)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.STRING)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.JSON)); break; case BOOLEAN: futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.DOUBLE)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.LONG)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.STRING)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.JSON)); break; case DOUBLE: futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.BOOLEAN)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.LONG)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.STRING)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.JSON)); break; case STRING: futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.BOOLEAN)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.DOUBLE)); futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.LONG)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.JSON)); + break; + case JSON: + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.BOOLEAN)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.DOUBLE)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.LONG)); + futures.add(saveNull(tenantId, entityId, tsKvEntry, ttl, partition, DataType.STRING)); break; } } @@ -411,6 +423,13 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem .set(5, tsKvEntry.getStrValue().orElse(null), String.class) .set(6, tsKvEntry.getLongValue().orElse(null), Long.class) .set(7, tsKvEntry.getDoubleValue().orElse(null), Double.class); + Optional jsonV = tsKvEntry.getJsonValue(); + if (jsonV.isPresent()) { + stmt.setString(8, tsKvEntry.getJsonValue().get()); + } else { + stmt.setToNull(8); + } + return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); } @@ -669,7 +688,12 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem if (boolV != null) { kvEntry = new BooleanDataEntry(key, boolV); } else { - log.warn("All values in key-value row are nullable "); + String jsonV = row.get(ModelConstants.JSON_VALUE_COLUMN, String.class); + if (StringUtils.isNoneEmpty(jsonV)) { + kvEntry = new JsonDataEntry(key, jsonV); + } else { + log.warn("All values in key-value row are nullable "); + } } } } @@ -772,8 +796,9 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem "," + ModelConstants.BOOLEAN_VALUE_COLUMN + "," + ModelConstants.STRING_VALUE_COLUMN + "," + ModelConstants.LONG_VALUE_COLUMN + - "," + ModelConstants.DOUBLE_VALUE_COLUMN + ")" + - " VALUES(?, ?, ?, ?, ?, ?, ?, ?)"); + "," + ModelConstants.DOUBLE_VALUE_COLUMN + + "," + ModelConstants.JSON_VALUE_COLUMN + ")" + + " VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?)"); } return latestInsertStmt; } @@ -812,7 +837,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem ModelConstants.STRING_VALUE_COLUMN + "," + ModelConstants.BOOLEAN_VALUE_COLUMN + "," + ModelConstants.LONG_VALUE_COLUMN + "," + - ModelConstants.DOUBLE_VALUE_COLUMN + " " + + ModelConstants.DOUBLE_VALUE_COLUMN + "," + + ModelConstants.JSON_VALUE_COLUMN + " " + "FROM " + ModelConstants.TS_KV_LATEST_CF + " " + "WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + @@ -829,7 +855,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem ModelConstants.STRING_VALUE_COLUMN + "," + ModelConstants.BOOLEAN_VALUE_COLUMN + "," + ModelConstants.LONG_VALUE_COLUMN + "," + - ModelConstants.DOUBLE_VALUE_COLUMN + " " + + ModelConstants.DOUBLE_VALUE_COLUMN + "," + + ModelConstants.JSON_VALUE_COLUMN + " " + "FROM " + ModelConstants.TS_KV_LATEST_CF + " " + "WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM); @@ -847,6 +874,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return ModelConstants.LONG_VALUE_COLUMN; case DOUBLE: return ModelConstants.DOUBLE_VALUE_COLUMN; + case JSON: + return ModelConstants.JSON_VALUE_COLUMN; default: throw new RuntimeException("Not implemented!"); } @@ -856,27 +885,23 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem switch (kvEntry.getDataType()) { case BOOLEAN: Optional booleanValue = kvEntry.getBooleanValue(); - if (booleanValue.isPresent()) { - stmt.setBool(column, booleanValue.get().booleanValue()); - } + booleanValue.ifPresent(b -> stmt.setBool(column, b)); break; case STRING: Optional stringValue = kvEntry.getStrValue(); - if (stringValue.isPresent()) { - stmt.setString(column, stringValue.get()); - } + stringValue.ifPresent(s -> stmt.setString(column, s)); break; case LONG: Optional longValue = kvEntry.getLongValue(); - if (longValue.isPresent()) { - stmt.setLong(column, longValue.get().longValue()); - } + longValue.ifPresent(l -> stmt.setLong(column, l)); break; case DOUBLE: Optional doubleValue = kvEntry.getDoubleValue(); - if (doubleValue.isPresent()) { - stmt.setDouble(column, doubleValue.get().doubleValue()); - } + doubleValue.ifPresent(d -> stmt.setDouble(column, d)); + break; + case JSON: + Optional jsonValue = kvEntry.getJsonValue(); + jsonValue.ifPresent(jsonObject -> stmt.setString(column, jsonObject)); break; } } diff --git a/dao/src/main/resources/cassandra/schema-entities.cql b/dao/src/main/resources/cassandra/schema-entities.cql index e9844f7b1c..de2b088cef 100644 --- a/dao/src/main/resources/cassandra/schema-entities.cql +++ b/dao/src/main/resources/cassandra/schema-entities.cql @@ -410,6 +410,7 @@ CREATE TABLE IF NOT EXISTS thingsboard.attributes_kv_cf ( str_v text, long_v bigint, dbl_v double, + json_v text, last_update_ts bigint, PRIMARY KEY ((entity_type, entity_id, attribute_type), attribute_key) ) WITH compaction = { 'class' : 'LeveledCompactionStrategy' }; diff --git a/dao/src/main/resources/cassandra/schema-ts.cql b/dao/src/main/resources/cassandra/schema-ts.cql index 338b420436..c0f4b74467 100644 --- a/dao/src/main/resources/cassandra/schema-ts.cql +++ b/dao/src/main/resources/cassandra/schema-ts.cql @@ -30,6 +30,7 @@ CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_cf ( str_v text, long_v bigint, dbl_v double, + json_v text, PRIMARY KEY (( entity_type, entity_id, key, partition ), ts) ); @@ -51,5 +52,6 @@ CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_latest_cf ( str_v text, long_v bigint, dbl_v double, + json_v text, PRIMARY KEY (( entity_type, entity_id ), key) ) WITH compaction = { 'class' : 'LeveledCompactionStrategy' }; diff --git a/dao/src/main/resources/sql/schema-entities-hsql.sql b/dao/src/main/resources/sql/schema-entities-hsql.sql new file mode 100644 index 0000000000..758aaafb10 --- /dev/null +++ b/dao/src/main/resources/sql/schema-entities-hsql.sql @@ -0,0 +1,251 @@ +-- +-- 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. +-- + + +CREATE TABLE IF NOT EXISTS admin_settings ( + id varchar(31) NOT NULL CONSTRAINT admin_settings_pkey PRIMARY KEY, + json_value varchar, + key varchar(255) +); + +CREATE TABLE IF NOT EXISTS alarm ( + id varchar(31) NOT NULL CONSTRAINT alarm_pkey PRIMARY KEY, + ack_ts bigint, + clear_ts bigint, + additional_info varchar, + end_ts bigint, + originator_id varchar(31), + originator_type integer, + propagate boolean, + severity varchar(255), + start_ts bigint, + status varchar(255), + tenant_id varchar(31), + propagate_relation_types varchar, + type varchar(255) +); + +CREATE TABLE IF NOT EXISTS asset ( + id varchar(31) NOT NULL CONSTRAINT asset_pkey PRIMARY KEY, + additional_info varchar, + customer_id varchar(31), + name varchar(255), + label varchar(255), + search_text varchar(255), + tenant_id varchar(31), + type varchar(255), + CONSTRAINT asset_name_unq_key UNIQUE (tenant_id, name) +); + +CREATE TABLE IF NOT EXISTS audit_log ( + id varchar(31) NOT NULL CONSTRAINT audit_log_pkey PRIMARY KEY, + tenant_id varchar(31), + customer_id varchar(31), + entity_id varchar(31), + entity_type varchar(255), + entity_name varchar(255), + user_id varchar(31), + user_name varchar(255), + action_type varchar(255), + action_data varchar(1000000), + action_status varchar(255), + action_failure_details varchar(1000000) +); + +CREATE TABLE IF NOT EXISTS attribute_kv ( + entity_type varchar(255), + entity_id varchar(31), + attribute_type varchar(255), + attribute_key varchar(255), + bool_v boolean, + str_v varchar(10000000), + long_v bigint, + dbl_v double precision, + json_v varchar(10000000), + last_update_ts bigint, + CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_type, entity_id, attribute_type, attribute_key) +); + +CREATE TABLE IF NOT EXISTS component_descriptor ( + id varchar(31) NOT NULL CONSTRAINT component_descriptor_pkey PRIMARY KEY, + actions varchar(255), + clazz varchar UNIQUE, + configuration_descriptor varchar, + name varchar(255), + scope varchar(255), + search_text varchar(255), + type varchar(255) +); + +CREATE TABLE IF NOT EXISTS customer ( + id varchar(31) NOT NULL CONSTRAINT customer_pkey PRIMARY KEY, + additional_info varchar, + address varchar, + address2 varchar, + city varchar(255), + country varchar(255), + email varchar(255), + phone varchar(255), + search_text varchar(255), + state varchar(255), + tenant_id varchar(31), + title varchar(255), + zip varchar(255) +); + +CREATE TABLE IF NOT EXISTS dashboard ( + id varchar(31) NOT NULL CONSTRAINT dashboard_pkey PRIMARY KEY, + configuration varchar(10000000), + assigned_customers varchar(1000000), + search_text varchar(255), + tenant_id varchar(31), + title varchar(255) +); + +CREATE TABLE IF NOT EXISTS device ( + id varchar(31) NOT NULL CONSTRAINT device_pkey PRIMARY KEY, + additional_info varchar, + customer_id varchar(31), + type varchar(255), + name varchar(255), + label varchar(255), + search_text varchar(255), + tenant_id varchar(31), + CONSTRAINT device_name_unq_key UNIQUE (tenant_id, name) +); + +CREATE TABLE IF NOT EXISTS device_credentials ( + id varchar(31) NOT NULL CONSTRAINT device_credentials_pkey PRIMARY KEY, + credentials_id varchar, + credentials_type varchar(255), + credentials_value varchar, + device_id varchar(31), + CONSTRAINT device_credentials_id_unq_key UNIQUE (credentials_id) +); + +CREATE TABLE IF NOT EXISTS event ( + id varchar(31) NOT NULL CONSTRAINT event_pkey PRIMARY KEY, + body varchar(10000000), + entity_id varchar(31), + entity_type varchar(255), + event_type varchar(255), + event_uid varchar(255), + tenant_id varchar(31), + CONSTRAINT event_unq_key UNIQUE (tenant_id, entity_type, entity_id, event_type, event_uid) +); + +CREATE TABLE IF NOT EXISTS relation ( + from_id varchar(31), + from_type varchar(255), + to_id varchar(31), + to_type varchar(255), + relation_type_group varchar(255), + relation_type varchar(255), + additional_info varchar, + CONSTRAINT relation_pkey PRIMARY KEY (from_id, from_type, relation_type_group, relation_type, to_id, to_type) +); + +CREATE TABLE IF NOT EXISTS tb_user ( + id varchar(31) NOT NULL CONSTRAINT tb_user_pkey PRIMARY KEY, + additional_info varchar, + authority varchar(255), + customer_id varchar(31), + email varchar(255) UNIQUE, + first_name varchar(255), + last_name varchar(255), + search_text varchar(255), + tenant_id varchar(31) +); + +CREATE TABLE IF NOT EXISTS tenant ( + id varchar(31) NOT NULL CONSTRAINT tenant_pkey PRIMARY KEY, + additional_info varchar, + address varchar, + address2 varchar, + city varchar(255), + country varchar(255), + email varchar(255), + phone varchar(255), + region varchar(255), + search_text varchar(255), + state varchar(255), + title varchar(255), + zip varchar(255) +); + +CREATE TABLE IF NOT EXISTS user_credentials ( + id varchar(31) NOT NULL CONSTRAINT user_credentials_pkey PRIMARY KEY, + activate_token varchar(255) UNIQUE, + enabled boolean, + password varchar(255), + reset_token varchar(255) UNIQUE, + user_id varchar(31) UNIQUE +); + +CREATE TABLE IF NOT EXISTS widget_type ( + id varchar(31) NOT NULL CONSTRAINT widget_type_pkey PRIMARY KEY, + alias varchar(255), + bundle_alias varchar(255), + descriptor varchar(1000000), + name varchar(255), + tenant_id varchar(31) +); + +CREATE TABLE IF NOT EXISTS widgets_bundle ( + id varchar(31) NOT NULL CONSTRAINT widgets_bundle_pkey PRIMARY KEY, + alias varchar(255), + search_text varchar(255), + tenant_id varchar(31), + title varchar(255) +); + +CREATE TABLE IF NOT EXISTS rule_chain ( + id varchar(31) NOT NULL CONSTRAINT rule_chain_pkey PRIMARY KEY, + additional_info varchar, + configuration varchar(10000000), + name varchar(255), + first_rule_node_id varchar(31), + root boolean, + debug_mode boolean, + search_text varchar(255), + tenant_id varchar(31) +); + +CREATE TABLE IF NOT EXISTS rule_node ( + id varchar(31) NOT NULL CONSTRAINT rule_node_pkey PRIMARY KEY, + rule_chain_id varchar(31), + additional_info varchar, + configuration varchar(10000000), + type varchar(255), + name varchar(255), + debug_mode boolean, + search_text varchar(255) +); + +CREATE TABLE IF NOT EXISTS entity_view ( + id varchar(31) NOT NULL CONSTRAINT entity_view_pkey PRIMARY KEY, + entity_id varchar(31), + entity_type varchar(255), + tenant_id varchar(31), + customer_id varchar(31), + type varchar(255), + name varchar(255), + keys varchar(10000000), + start_ts bigint, + end_ts bigint, + search_text varchar(255), + additional_info varchar +); diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index f59b1045bc..55893fc124 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -74,6 +74,7 @@ CREATE TABLE IF NOT EXISTS attribute_kv ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v json, last_update_ts bigint, CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_type, entity_id, attribute_type, attribute_key) ); diff --git a/dao/src/main/resources/sql/schema-timescale.sql b/dao/src/main/resources/sql/schema-timescale.sql index e8cf0de263..7251d8be4e 100644 --- a/dao/src/main/resources/sql/schema-timescale.sql +++ b/dao/src/main/resources/sql/schema-timescale.sql @@ -25,6 +25,7 @@ CREATE TABLE IF NOT EXISTS tenant_ts_kv ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v json, CONSTRAINT tenant_ts_kv_pkey PRIMARY KEY (tenant_id, entity_id, key, ts) ); @@ -42,5 +43,6 @@ CREATE TABLE IF NOT EXISTS ts_kv_latest ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v json, CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key) ); \ No newline at end of file diff --git a/dao/src/main/resources/sql/schema-ts-hsql.sql b/dao/src/main/resources/sql/schema-ts-hsql.sql index c29d7e2ed7..eb053a7a84 100644 --- a/dao/src/main/resources/sql/schema-ts-hsql.sql +++ b/dao/src/main/resources/sql/schema-ts-hsql.sql @@ -24,6 +24,7 @@ CREATE TABLE IF NOT EXISTS ts_kv ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v varchar(10000000), CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_id, key, ts) ); @@ -35,6 +36,7 @@ CREATE TABLE IF NOT EXISTS ts_kv_latest ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v varchar(10000000), CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key) ); diff --git a/dao/src/main/resources/sql/schema-ts-psql.sql b/dao/src/main/resources/sql/schema-ts-psql.sql index 465c2d51e3..32b6762c8e 100644 --- a/dao/src/main/resources/sql/schema-ts-psql.sql +++ b/dao/src/main/resources/sql/schema-ts-psql.sql @@ -21,7 +21,8 @@ CREATE TABLE IF NOT EXISTS ts_kv ( bool_v boolean, str_v varchar(10000000), long_v bigint, - dbl_v double precision + dbl_v double precision, + json_v json ) PARTITION BY RANGE (ts); CREATE TABLE IF NOT EXISTS ts_kv_latest ( @@ -32,6 +33,7 @@ CREATE TABLE IF NOT EXISTS ts_kv_latest ( str_v varchar(10000000), long_v bigint, dbl_v double precision, + json_v json, CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key) ); diff --git a/dao/src/test/java/org/thingsboard/server/dao/JpaDaoTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/JpaDaoTestSuite.java index f4d11b328e..0af8603bf1 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/JpaDaoTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/JpaDaoTestSuite.java @@ -30,7 +30,7 @@ public class JpaDaoTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/system-data.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/system-data.sql"), "sql/drop-all-tables.sql", "sql-test.properties" ); diff --git a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java index caddbabc35..7ebab237a8 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java @@ -30,7 +30,7 @@ public class SqlDaoServiceTestSuite { @ClassRule public static CustomSqlUnit sqlUnit = new CustomSqlUnit( - Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities.sql", "sql/schema-entities-idx.sql", "sql/system-data.sql", "sql/system-test.sql"), + Arrays.asList("sql/schema-ts-hsql.sql", "sql/schema-entities-hsql.sql", "sql/schema-entities-idx.sql", "sql/system-data.sql", "sql/system-test.sql"), "sql/drop-all-tables.sql", "sql-test.properties" ); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index 4db628a7f4..be4e6b4e30 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.gson.JsonParseException; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.BooleanUtils; import org.thingsboard.rule.engine.api.TbContext; @@ -33,6 +34,7 @@ import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; +import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ConcurrentHashMap; @@ -77,7 +79,8 @@ public abstract class TbAbstractGetAttributesNode findEntityIdAsync(TbContext ctx, TbMsg msg); @@ -168,6 +171,12 @@ public abstract class TbAbstractGetAttributesNode