Browse Source

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>
pull/2419/head
Andrew Shvayka 7 years ago
committed by GitHub
parent
commit
03f5375a02
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 7
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 2
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  3. 43
      application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
  4. 26
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  5. 1
      application/src/main/proto/cluster.proto
  6. 2
      application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java
  7. 2
      application/src/test/java/org/thingsboard/server/mqtt/MqttSqlTestSuite.java
  8. 2
      application/src/test/java/org/thingsboard/server/rules/RuleEngineSqlTestSuite.java
  9. 2
      application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java
  10. 3
      common/data/src/main/java/org/thingsboard/server/common/data/SearchTextBasedWithAdditionalInfo.java
  11. 7
      common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseAttributeKvEntry.java
  12. 5
      common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicKvEntry.java
  13. 5
      common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java
  14. 2
      common/data/src/main/java/org/thingsboard/server/common/data/kv/DataType.java
  15. 69
      common/data/src/main/java/org/thingsboard/server/common/data/kv/JsonDataEntry.java
  16. 2
      common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java
  17. 40
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java
  18. 2
      common/transport/transport-api/src/main/proto/transport.proto
  19. 40
      dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
  20. 17
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  21. 8
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java
  22. 8
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AttributeKvEntity.java
  23. 18
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java
  24. 6
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestEntity.java
  25. 15
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/psql/TsKvEntity.java
  26. 17
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java
  27. 118
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvInsertRepository.java
  28. 34
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/HsqlAttributesInsertRepository.java
  29. 1
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  30. 19
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/PsqlAttributesInsertRepository.java
  31. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  32. 12
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java
  33. 3
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java
  34. 8
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java
  35. 33
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java
  36. 3
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java
  37. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java
  38. 1
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java
  39. 21
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java
  40. 3
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java
  41. 9
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java
  42. 19
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java
  43. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  44. 48
      dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java
  45. 59
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  46. 1
      dao/src/main/resources/cassandra/schema-entities.cql
  47. 2
      dao/src/main/resources/cassandra/schema-ts.cql
  48. 251
      dao/src/main/resources/sql/schema-entities-hsql.sql
  49. 1
      dao/src/main/resources/sql/schema-entities.sql
  50. 2
      dao/src/main/resources/sql/schema-timescale.sql
  51. 2
      dao/src/main/resources/sql/schema-ts-hsql.sql
  52. 4
      dao/src/main/resources/sql/schema-ts-psql.sql
  53. 2
      dao/src/test/java/org/thingsboard/server/dao/JpaDaoTestSuite.java
  54. 2
      dao/src/test/java/org/thingsboard/server/dao/SqlDaoServiceTestSuite.java
  55. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  56. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

7
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();
}

2
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());
}

43
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<List<AttributeKvEntry>>() {
@Override
public void onSuccess(List<AttributeKvEntry> attributes) {
List<AttributeData> values = attributes.stream().map(attribute -> new AttributeData(attribute.getLastUpdateTs(),
attribute.getKey(), attribute.getValue())).collect(Collectors.toList());
List<AttributeData> 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();
}
}

26
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<Double> doubleValue = attr.getDoubleValue();
doubleValue.ifPresent(dataBuilder::setDoubleValue);
break;
case JSON:
Optional<String> jsonValue = attr.getJsonValue();
jsonValue.ifPresent(dataBuilder::setJsonValue);
break;
case STRING:
Optional<String> 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;
}

1
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 {

2
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");
}

2
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");
}

2
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");
}

2
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");

3
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<I extends UUIDBased> extends SearchTextBased<I> 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<I extends UUIDBased> ext
public static void setJson(JsonNode json, Consumer<JsonNode> jsonConsumer, Consumer<byte[]> 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);
}

7
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<String> getJsonValue() {
return kv.getJsonValue();
}
@Override
public String getValueAsString() {
return kv.getValueAsString();

5
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<String> getJsonValue() {
return Optional.ofNullable(null);
}
@Override
public boolean equals(Object o) {
if (this == o) return true;

5
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<String> getJsonValue() {
return kv.getJsonValue();
}
@Override
public Object getValue() {
return kv.getValue();

2
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;
}

69
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<String> 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;
}
}

2
common/data/src/main/java/org/thingsboard/server/common/data/kv/KvEntry.java

@ -37,6 +37,8 @@ public interface KvEntry extends Serializable {
Optional<Double> getDoubleValue();
Optional<String> getJsonValue();
String getValueAsString();
Object getValue();

40
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<TsKvProto> 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<AttributeKvEntry> 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);
}

2
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 {

40
dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java

@ -112,31 +112,18 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem
@Override
public ListenableFuture<Void> 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<Boolean> booleanValue = attribute.getBooleanValue();
if (booleanValue.isPresent()) {
stmt.setBool(6, booleanValue.get());
} else {
stmt.setToNull(6);
}
Optional<Long> longValue = attribute.getLongValue();
if (longValue.isPresent()) {
stmt.setLong(7, longValue.get());
} else {
stmt.setToNull(7);
}
Optional<Double> 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;
}

17
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) {

8
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<TsKvEntry> {
@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<TsKvEntry> {
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);
}

8
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<AttributeKvEntry>, 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<AttributeKvEntry>, 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);
}
}

18
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<TsKvE
}
}
public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount) {
public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount) {
if (!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) {
if (booleanValueCount != 0) {
this.longValue = booleanValueCount;
} else if (strValueCount != 0) {
this.longValue = strValueCount;
} else if (jsonValueCount != 0) {
this.longValue = jsonValueCount;
} else {
this.longValue = longValueCount + doubleValueCount;
}

6
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestEntity.java

@ -52,6 +52,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN;
@ColumnResult(name = "boolValue", type = Boolean.class),
@ColumnResult(name = "longValue", type = Long.class),
@ColumnResult(name = "doubleValue", type = Double.class),
@ColumnResult(name = "jsonValue", type = String.class),
@ColumnResult(name = "ts", type = Long.class),
}
@ -74,13 +75,13 @@ public final class TsKvLatestEntity extends AbstractTsKvEntity {
@Override
public boolean isNotEmpty() {
return strValue != null || longValue != null || doubleValue != null || booleanValue != null;
return strValue != null || longValue != null || doubleValue != null || booleanValue != null || jsonValue != null;
}
public TsKvLatestEntity() {
}
public TsKvLatestEntity(UUID entityId, Integer key, String strKey, String strValue, Boolean boolValue, Long longValue, Double doubleValue, Long ts) {
public TsKvLatestEntity(UUID entityId, Integer key, String strKey, String strValue, Boolean boolValue, Long longValue, Double doubleValue, String jsonValue, Long ts) {
this.entityId = entityId;
this.key = key;
this.ts = ts;
@ -88,6 +89,7 @@ public final class TsKvLatestEntity extends AbstractTsKvEntity {
this.doubleValue = doubleValue;
this.strValue = strValue;
this.booleanValue = boolValue;
this.jsonValue = jsonValue;
this.strKey = strKey;
}
}

15
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/psql/TsKvEntity.java

@ -16,14 +16,6 @@
package org.thingsboard.server.dao.model.sqlts.psql;
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.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;
@ -31,10 +23,7 @@ import javax.persistence.Entity;
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.KEY_COLUMN;
@Data
@ -93,12 +82,14 @@ public final class TsKvEntity extends AbstractTsKvEntity {
}
}
public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount) {
public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount) {
if (!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) {
if (booleanValueCount != 0) {
this.longValue = booleanValueCount;
} else if (strValueCount != 0) {
this.longValue = strValueCount;
} else if (jsonValueCount != 0) {
this.longValue = jsonValueCount;
} else {
this.longValue = longValueCount + doubleValueCount;
}

17
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java

@ -18,12 +18,6 @@ package org.thingsboard.server.dao.model.sqlts.timescale;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.springframework.util.StringUtils;
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;
@ -39,10 +33,8 @@ import javax.persistence.NamedNativeQuery;
import javax.persistence.SqlResultSetMapping;
import javax.persistence.SqlResultSetMappings;
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.KEY_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.TENANT_ID_COLUMN;
import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_AVG;
@ -92,6 +84,7 @@ import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.F
@ColumnResult(name = "strValueCount", type = Long.class),
@ColumnResult(name = "longValueCount", type = Long.class),
@ColumnResult(name = "doubleValueCount", type = Long.class),
@ColumnResult(name = "jsonValueCount", type = Long.class),
}
)
}),
@ -179,13 +172,15 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity implements ToD
}
}
public TimescaleTsKvEntity(Long tsBucket, Long interval, Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount) {
if (!isAllNull(tsBucket, interval, booleanValueCount, strValueCount, longValueCount, doubleValueCount)) {
public TimescaleTsKvEntity(Long tsBucket, Long interval, Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount) {
if (!isAllNull(tsBucket, interval, booleanValueCount, strValueCount, longValueCount, doubleValueCount, jsonValueCount)) {
this.ts = tsBucket + interval / 2;
if (booleanValueCount != 0) {
this.longValue = booleanValueCount;
} else if (strValueCount != 0) {
this.longValue = strValueCount;
} else if (jsonValueCount != 0) {
this.longValue = jsonValueCount;
} else {
this.longValue = longValueCount + doubleValueCount;
}
@ -194,6 +189,6 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity implements ToD
@Override
public boolean isNotEmpty() {
return ts != null && (strValue != null || longValue != null || doubleValue != null || booleanValue != null);
return ts != null && (strValue != null || longValue != null || doubleValue != null || booleanValue != null || jsonValue != null);
}
}

118
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvInsertRepository.java

@ -18,7 +18,6 @@ package org.thingsboard.server.dao.sql.attributes;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository;
@ -28,8 +27,6 @@ import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.model.sql.AttributeKvEntity;
import org.thingsboard.server.dao.util.SqlDao;
import javax.persistence.EntityManager;
import javax.persistence.PersistenceContext;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.Types;
@ -45,19 +42,14 @@ public abstract class AttributeKvInsertRepository {
private static final ThreadLocal<Pattern> 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<AttributeKvEntity> 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

34
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<AttributeKvEntity> 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());
});
});
}

1
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);
}

19
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;
}
}

2
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);
}

12
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<TsKvEntity> {
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<EntityContainer<TsKvEntity>> entities) {
@ -76,6 +76,8 @@ public class HsqlInsertTsRepository extends AbstractInsertRepository implements
} else {
ps.setNull(7, Types.DOUBLE);
}
ps.setString(8, tsKvEntity.getJsonValue());
}
@Override

3
dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java

@ -97,7 +97,8 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
@Query("SELECT new TsKvEntity(SUM(CASE WHEN tskv.booleanValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.strValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.jsonValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findCount(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,

8
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

33
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<TsKvLatestEntity> 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

3
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

4
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<TsKvLatestEntity, TsKvLatestCompositeKey> {
List<TsKvLatestEntity> findAllByEntityId(UUID entityId);
}

1
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()));

21
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<EntityContainer<TsKvEntity>> 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()));
}
}

3
dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/TsKvPsqlRepository.java

@ -98,7 +98,8 @@ public interface TsKvPsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
@Query("SELECT new TsKvEntity(SUM(CASE WHEN tskv.booleanValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.strValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.jsonValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findCount(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,

9
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;

19
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<TimescaleTsKvEntity> {
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<EntityContainer<TimescaleTsKvEntity>> 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

2
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));
}

48
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<TsKvEntry> 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;

59
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<String> 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<Boolean> booleanValue = kvEntry.getBooleanValue();
if (booleanValue.isPresent()) {
stmt.setBool(column, booleanValue.get().booleanValue());
}
booleanValue.ifPresent(b -> stmt.setBool(column, b));
break;
case STRING:
Optional<String> stringValue = kvEntry.getStrValue();
if (stringValue.isPresent()) {
stmt.setString(column, stringValue.get());
}
stringValue.ifPresent(s -> stmt.setString(column, s));
break;
case LONG:
Optional<Long> longValue = kvEntry.getLongValue();
if (longValue.isPresent()) {
stmt.setLong(column, longValue.get().longValue());
}
longValue.ifPresent(l -> stmt.setLong(column, l));
break;
case DOUBLE:
Optional<Double> doubleValue = kvEntry.getDoubleValue();
if (doubleValue.isPresent()) {
stmt.setDouble(column, doubleValue.get().doubleValue());
}
doubleValue.ifPresent(d -> stmt.setDouble(column, d));
break;
case JSON:
Optional<String> jsonValue = kvEntry.getJsonValue();
jsonValue.ifPresent(jsonObject -> stmt.setString(column, jsonObject));
break;
}
}

1
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' };

2
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' };

251
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
);

1
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)
);

2
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)
);

2
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)
);

4
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)
);

2
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"
);

2
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"
);

11
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<C extends TbGetAttributesNodeC
}
@Override
public void destroy() { }
public void destroy() {
}
protected abstract ListenableFuture<T> findEntityIdAsync(TbContext ctx, TbMsg msg);
@ -168,6 +171,12 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
case DOUBLE:
value.put(VALUE, r.getDoubleValue().get());
break;
case JSON:
try {
value.set(VALUE, mapper.readTree(r.getJsonValue().get()));
} catch (IOException e) {
throw new JsonParseException("Can't parse jsonValue: " + r.getJsonValue().get(), e);
}
}
msg.getMetaData().putValue(r.getKey(), value.toString());
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.gson.JsonParseException;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -39,6 +40,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
@ -180,6 +182,13 @@ public class TbGetTelemetryNode implements TbNode {
case DOUBLE:
obj.put("value", entry.getDoubleValue().get());
break;
case JSON:
try {
obj.set("value", mapper.readTree(entry.getJsonValue().get()));
} catch (IOException e) {
throw new JsonParseException("Can't parse jsonValue: " + entry.getJsonValue().get(), e);
}
break;
}
return obj;
}

Loading…
Cancel
Save