Browse Source

Refactor

pull/5023/head
Viacheslav Klimov 5 years ago
parent
commit
38f6e314b1
  1. 11
      application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java
  2. 33
      application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java
  3. 11
      application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java
  4. 80
      application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java
  5. 18
      application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java
  6. 55
      application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java
  7. 93
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java
  8. 6
      common/transport/transport-api/src/test/java/JsonConverterTest.java

11
application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java

@ -26,6 +26,7 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.action.EntityActionService;
import org.thingsboard.server.service.importing.AbstractBulkImportService;
import org.thingsboard.server.service.importing.BulkImportColumnType;
import org.thingsboard.server.service.importing.BulkImportRequest;
import org.thingsboard.server.service.importing.ImportedEntityInfo;
import org.thingsboard.server.service.security.AccessValidator;
@ -49,12 +50,12 @@ public class AssetBulkImportService extends AbstractBulkImportService<Asset> {
}
@Override
protected ImportedEntityInfo<Asset> saveEntity(BulkImportRequest importRequest, Map<BulkImportRequest.ColumnMapping, String> entityData, SecurityUser user) {
protected ImportedEntityInfo<Asset> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) {
ImportedEntityInfo<Asset> importedEntityInfo = new ImportedEntityInfo<>();
Asset asset = new Asset();
asset.setTenantId(user.getTenantId());
setAssetFields(asset, entityData);
setAssetFields(asset, fields);
Asset existingAsset = assetService.findAssetByTenantIdAndName(user.getTenantId(), asset.getName());
if (existingAsset != null && importRequest.getMapping().getUpdate()) {
@ -69,10 +70,10 @@ public class AssetBulkImportService extends AbstractBulkImportService<Asset> {
return importedEntityInfo;
}
private void setAssetFields(Asset asset, Map<BulkImportRequest.ColumnMapping, String> data) {
private void setAssetFields(Asset asset, Map<BulkImportColumnType, String> fields) {
ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(asset.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode);
data.forEach((columnMapping, value) -> {
switch (columnMapping.getType()) {
fields.forEach((columnType, value) -> {
switch (columnType) {
case NAME:
asset.setName(value);
break;

33
application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java

@ -58,7 +58,6 @@ import java.util.EnumSet;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@Service
@ -80,12 +79,12 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
}
@Override
protected ImportedEntityInfo<Device> saveEntity(BulkImportRequest importRequest, Map<BulkImportRequest.ColumnMapping, String> entityData, SecurityUser user) {
protected ImportedEntityInfo<Device> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) {
ImportedEntityInfo<Device> importedEntityInfo = new ImportedEntityInfo<>();
Device device = new Device();
device.setTenantId(user.getTenantId());
setDeviceFields(device, entityData);
setDeviceFields(device, fields);
Device existingDevice = deviceService.findDeviceByTenantIdAndName(user.getTenantId(), device.getName());
if (existingDevice != null && importRequest.getMapping().getUpdate()) {
@ -95,7 +94,7 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
device = existingDevice;
}
DeviceCredentials deviceCredentials = createDeviceCredentials(entityData);
DeviceCredentials deviceCredentials = createDeviceCredentials(fields);
if (deviceCredentials.getCredentialsType() != null) {
if (deviceCredentials.getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) {
setUpLwM2mDeviceProfile(user.getTenantId(), device);
@ -121,10 +120,10 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
return importedEntityInfo;
}
private void setDeviceFields(Device device, Map<BulkImportRequest.ColumnMapping, String> data) {
private void setDeviceFields(Device device, Map<BulkImportColumnType, String> fields) {
ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(device.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode);
data.forEach((columnMapping, value) -> {
switch (columnMapping.getType()) {
fields.forEach((columnType, value) -> {
switch (columnType) {
case NAME:
device.setName(value);
break;
@ -146,25 +145,25 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
}
@SneakyThrows
private DeviceCredentials createDeviceCredentials(Map<BulkImportRequest.ColumnMapping, String> data) {
Set<BulkImportColumnType> columns = data.keySet().stream().map(BulkImportRequest.ColumnMapping::getType).collect(Collectors.toSet());
private DeviceCredentials createDeviceCredentials(Map<BulkImportColumnType, String> fields) {
Set<BulkImportColumnType> columns = fields.keySet();
DeviceCredentials credentials = new DeviceCredentials();
if (columns.contains(BulkImportColumnType.ACCESS_TOKEN)) {
credentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN);
credentials.setCredentialsId(getByColumnType(BulkImportColumnType.ACCESS_TOKEN, data));
credentials.setCredentialsId(fields.get(BulkImportColumnType.ACCESS_TOKEN));
} else if (CollectionUtils.containsAny(columns, EnumSet.of(BulkImportColumnType.MQTT_CLIENT_ID, BulkImportColumnType.MQTT_USER_NAME, BulkImportColumnType.MQTT_PASSWORD))) {
credentials.setCredentialsType(DeviceCredentialsType.MQTT_BASIC);
BasicMqttCredentials basicMqttCredentials = new BasicMqttCredentials();
basicMqttCredentials.setClientId(getByColumnType(BulkImportColumnType.MQTT_CLIENT_ID, data));
basicMqttCredentials.setUserName(getByColumnType(BulkImportColumnType.MQTT_USER_NAME, data));
basicMqttCredentials.setPassword(getByColumnType(BulkImportColumnType.MQTT_PASSWORD, data));
basicMqttCredentials.setClientId(fields.get(BulkImportColumnType.MQTT_CLIENT_ID));
basicMqttCredentials.setUserName(fields.get(BulkImportColumnType.MQTT_USER_NAME));
basicMqttCredentials.setPassword(fields.get(BulkImportColumnType.MQTT_PASSWORD));
credentials.setCredentialsValue(JacksonUtil.toString(basicMqttCredentials));
} else if (columns.contains(BulkImportColumnType.X509)) {
credentials.setCredentialsType(DeviceCredentialsType.X509_CERTIFICATE);
credentials.setCredentialsValue(getByColumnType(BulkImportColumnType.X509, data));
credentials.setCredentialsValue(fields.get(BulkImportColumnType.X509));
} else if (columns.contains(BulkImportColumnType.LWM2M_CLIENT_ENDPOINT)) {
credentials.setCredentialsType(DeviceCredentialsType.LWM2M_CREDENTIALS);
ObjectNode lwm2mCredentials = JacksonUtil.newObjectNode();
@ -173,7 +172,7 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
Stream.of(BulkImportColumnType.LWM2M_CLIENT_ENDPOINT, BulkImportColumnType.LWM2M_CLIENT_SECURITY_CONFIG_MODE,
BulkImportColumnType.LWM2M_CLIENT_IDENTITY, BulkImportColumnType.LWM2M_CLIENT_KEY, BulkImportColumnType.LWM2M_CLIENT_CERT)
.forEach(lwm2mClientProperty -> {
String value = getByColumnType(lwm2mClientProperty, data);
String value = fields.get(lwm2mClientProperty);
if (value != null) {
client.set(lwm2mClientProperty.getKey(), new TextNode(value));
}
@ -187,7 +186,7 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
Stream.of(BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_SECURITY_MODE, BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_PUBLIC_KEY_OR_ID,
BulkImportColumnType.LWM2M_BOOTSTRAP_SERVER_SECRET_KEY)
.forEach(lwm2mBootstrapServerProperty -> {
String value = getByColumnType(lwm2mBootstrapServerProperty, data);
String value = fields.get(lwm2mBootstrapServerProperty);
if (value != null) {
bootstrapServer.set(lwm2mBootstrapServerProperty.getKey(), new TextNode(value));
}
@ -197,7 +196,7 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
Stream.of(BulkImportColumnType.LWM2M_SERVER_SECURITY_MODE, BulkImportColumnType.LWM2M_SERVER_CLIENT_PUBLIC_KEY_OR_ID,
BulkImportColumnType.LWM2M_SERVER_CLIENT_SECRET_KEY)
.forEach(lwm2mServerProperty -> {
String value = getByColumnType(lwm2mServerProperty, data);
String value = fields.get(lwm2mServerProperty);
if (value != null) {
lwm2mServer.set(lwm2mServerProperty.getKey(), new TextNode(value));
}

11
application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java

@ -26,6 +26,7 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.action.EntityActionService;
import org.thingsboard.server.service.importing.AbstractBulkImportService;
import org.thingsboard.server.service.importing.BulkImportColumnType;
import org.thingsboard.server.service.importing.BulkImportRequest;
import org.thingsboard.server.service.importing.ImportedEntityInfo;
import org.thingsboard.server.service.security.AccessValidator;
@ -49,12 +50,12 @@ public class EdgeBulkImportService extends AbstractBulkImportService<Edge> {
}
@Override
protected ImportedEntityInfo<Edge> saveEntity(BulkImportRequest importRequest, Map<BulkImportRequest.ColumnMapping, String> entityData, SecurityUser user) {
protected ImportedEntityInfo<Edge> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user) {
ImportedEntityInfo<Edge> importedEntityInfo = new ImportedEntityInfo<>();
Edge edge = new Edge();
edge.setTenantId(user.getTenantId());
setEdgeFields(edge, entityData);
setEdgeFields(edge, fields);
Edge existingEdge = edgeService.findEdgeByTenantIdAndName(user.getTenantId(), edge.getName());
if (existingEdge != null && importRequest.getMapping().getUpdate()) {
@ -69,10 +70,10 @@ public class EdgeBulkImportService extends AbstractBulkImportService<Edge> {
return importedEntityInfo;
}
private void setEdgeFields(Edge edge, Map<BulkImportRequest.ColumnMapping, String> data) {
private void setEdgeFields(Edge edge, Map<BulkImportColumnType, String> fields) {
ObjectNode additionalInfo = (ObjectNode) Optional.ofNullable(edge.getAdditionalInfo()).orElseGet(JacksonUtil::newObjectNode);
data.forEach((columnMapping, value) -> {
switch (columnMapping.getType()) {
fields.forEach((columnType, value) -> {
switch (columnType) {
case NAME:
edge.setName(value);
break;

80
application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.importing;
import com.google.common.util.concurrent.FutureCallback;
import com.google.gson.JsonObject;
import com.google.gson.JsonPrimitive;
import lombok.Data;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import org.apache.commons.lang3.StringUtils;
@ -29,6 +30,7 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.DataType;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
@ -42,9 +44,12 @@ import org.thingsboard.server.service.security.permission.AccessControlService;
import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import org.thingsboard.server.utils.CsvUtils;
import org.thingsboard.server.utils.TypeCastUtil;
import javax.annotation.Nullable;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
@ -73,12 +78,12 @@ public abstract class AbstractBulkImportService<E extends BaseData<? extends Ent
parseData(request).forEach(entityData -> {
i.incrementAndGet();
try {
ImportedEntityInfo<E> importedEntityInfo = saveEntity(request, entityData, user);
ImportedEntityInfo<E> importedEntityInfo = saveEntity(request, entityData.getFields(), user);
onEntityImported.accept(importedEntityInfo);
E entity = importedEntityInfo.getEntity();
saveKvs(user, entity, entityData);
saveKvs(user, entity, entityData.getKvs());
if (importedEntityInfo.getRelatedError() != null) {
throw new RuntimeException(importedEntityInfo.getRelatedError());
@ -98,20 +103,21 @@ public abstract class AbstractBulkImportService<E extends BaseData<? extends Ent
return result;
}
protected abstract ImportedEntityInfo<E> saveEntity(BulkImportRequest importRequest, Map<ColumnMapping, String> entityData, SecurityUser user);
protected abstract ImportedEntityInfo<E> saveEntity(BulkImportRequest importRequest, Map<BulkImportColumnType, String> fields, SecurityUser user);
/*
* Attributes' values are firstly added to JsonObject in order to then make some type cast,
* because we get all values as strings from CSV
* */
private void saveKvs(SecurityUser user, E entity, Map<ColumnMapping, String> data) {
Stream.of(BulkImportColumnType.SHARED_ATTRIBUTE, BulkImportColumnType.SERVER_ATTRIBUTE, BulkImportColumnType.TIMESERIES)
private void saveKvs(SecurityUser user, E entity, Map<ColumnMapping, ParsedValue> data) {
Arrays.stream(BulkImportColumnType.values())
.filter(BulkImportColumnType::isKv)
.map(kvType -> {
JsonObject kvs = new JsonObject();
data.entrySet().stream()
.filter(dataEntry -> dataEntry.getKey().getType() == kvType &&
StringUtils.isNotEmpty(dataEntry.getKey().getKey()))
.forEach(dataEntry -> kvs.add(dataEntry.getKey().getKey(), new JsonPrimitive(dataEntry.getValue())));
.forEach(dataEntry -> kvs.add(dataEntry.getKey().getKey(), dataEntry.getValue().toJsonPrimitive()));
return Map.entry(kvType, kvs);
})
.filter(kvsEntry -> kvsEntry.getValue().entrySet().size() > 0)
@ -127,7 +133,7 @@ public abstract class AbstractBulkImportService<E extends BaseData<? extends Ent
@SneakyThrows
private void saveTelemetry(SecurityUser user, E entity, Map.Entry<BulkImportColumnType, JsonObject> kvsEntry) {
List<TsKvEntry> timeseries = JsonConverter.convertToTelemetry(kvsEntry.getValue(), System.currentTimeMillis(), false, true)
List<TsKvEntry> timeseries = JsonConverter.convertToTelemetry(kvsEntry.getValue(), System.currentTimeMillis())
.entrySet().stream()
.flatMap(entry -> entry.getValue().stream().map(kvEntry -> new BasicTsKvEntry(entry.getKey(), kvEntry)))
.collect(Collectors.toList());
@ -155,7 +161,7 @@ public abstract class AbstractBulkImportService<E extends BaseData<? extends Ent
@SneakyThrows
private void saveAttributes(SecurityUser user, E entity, Map.Entry<BulkImportColumnType, JsonObject> kvsEntry, BulkImportColumnType kvType) {
String scope = kvType.getKey();
List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue(), true));
List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue()));
accessValidator.validateEntityAndCallback(user, Operation.WRITE_ATTRIBUTES, entity.getId(), (result, tenantId, entityId) -> {
tsSubscriptionService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<>() {
@ -178,24 +184,62 @@ public abstract class AbstractBulkImportService<E extends BaseData<? extends Ent
});
}
protected final String getByColumnType(BulkImportColumnType bulkImportColumnType, Map<ColumnMapping, String> data) {
return data.entrySet().stream().filter(entry -> entry.getKey().getType() == bulkImportColumnType).findFirst().map(Map.Entry::getValue).orElse(null);
}
private List<Map<ColumnMapping, String>> parseData(BulkImportRequest request) throws Exception {
private List<EntityData> parseData(BulkImportRequest request) throws Exception {
List<List<String>> records = CsvUtils.parseCsv(request.getFile(), request.getMapping().getDelimiter());
if (request.getMapping().getHeader()) {
records.remove(0);
}
List<ColumnMapping> columnsMappings = request.getMapping().getColumns();
return records.stream()
.map(record -> Stream.iterate(0, i -> i < record.size(), i -> i + 1)
.map(i -> Map.entry(columnsMappings.get(i), record.get(i)))
.filter(entry -> StringUtils.isNotEmpty(entry.getValue()))
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)))
.map(record -> {
EntityData entityData = new EntityData();
Stream.iterate(0, i -> i < record.size(), i -> i + 1)
.map(i -> Map.entry(columnsMappings.get(i), record.get(i)))
.filter(entry -> StringUtils.isNotEmpty(entry.getValue()))
.forEach(entry -> {
if (!entry.getKey().getType().isKv()) {
entityData.getFields().put(entry.getKey().getType(), entry.getValue());
} else {
Map.Entry<DataType, Object> castResult = TypeCastUtil.castValue(entry.getValue());
entityData.getKvs().put(entry.getKey(), new ParsedValue(castResult.getValue(), castResult.getKey()));
}
});
return entityData;
})
.collect(Collectors.toList());
}
@Data
protected static class EntityData {
private final Map<BulkImportColumnType, String> fields = new HashMap<>();
private final Map<ColumnMapping, ParsedValue> kvs = new HashMap<>();
}
@Data
protected static class ParsedValue {
private final Object value;
private final DataType dataType;
public JsonPrimitive toJsonPrimitive() {
switch (dataType) {
case STRING:
return new JsonPrimitive((String) value);
case LONG:
return new JsonPrimitive((Long) value);
case DOUBLE:
return new JsonPrimitive((Double) value);
case BOOLEAN:
return new JsonPrimitive((Boolean) value);
default:
return null;
}
}
public String stringValue() {
return value.toString();
}
}
}

18
application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java

@ -18,13 +18,14 @@ package org.thingsboard.server.service.importing;
import lombok.Getter;
import org.thingsboard.server.common.data.DataConstants;
@Getter
public enum BulkImportColumnType {
NAME,
TYPE,
LABEL,
SHARED_ATTRIBUTE(DataConstants.SHARED_SCOPE),
SERVER_ATTRIBUTE(DataConstants.SERVER_SCOPE),
TIMESERIES,
SHARED_ATTRIBUTE(DataConstants.SHARED_SCOPE, true),
SERVER_ATTRIBUTE(DataConstants.SERVER_SCOPE, true),
TIMESERIES(true),
ACCESS_TOKEN,
X509,
MQTT_CLIENT_ID,
@ -48,8 +49,8 @@ public enum BulkImportColumnType {
ROUTING_KEY,
SECRET;
@Getter
private String key;
private boolean isKv = false;
BulkImportColumnType() {
}
@ -57,4 +58,13 @@ public enum BulkImportColumnType {
BulkImportColumnType(String key) {
this.key = key;
}
BulkImportColumnType(boolean isKv) {
this.isKv = isKv;
}
BulkImportColumnType(String key, boolean isKv) {
this.key = key;
this.isKv = isKv;
}
}

55
application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java

@ -0,0 +1,55 @@
/**
* Copyright © 2016-2021 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.utils;
import org.apache.commons.lang3.math.NumberUtils;
import org.thingsboard.server.common.data.kv.DataType;
import java.math.BigDecimal;
import java.util.Map;
public class TypeCastUtil {
private TypeCastUtil() {}
public static Map.Entry<DataType, Object> castValue(String value) {
if (isNumber(value)) {
String formattedValue = value.replace(',', '.');
try {
BigDecimal bd = new BigDecimal(formattedValue);
if (bd.stripTrailingZeros().scale() > 0 || isSimpleDouble(formattedValue)) {
if (bd.scale() <= 16) {
return Map.entry(DataType.DOUBLE, bd.doubleValue());
}
} else {
return Map.entry(DataType.LONG, bd.longValueExact());
}
} catch (RuntimeException ignored) {}
} else if (value.equalsIgnoreCase("true") || value.equalsIgnoreCase("false")) {
return Map.entry(DataType.BOOLEAN, Boolean.parseBoolean(value));
}
return Map.entry(DataType.STRING, value);
}
private static boolean isNumber(String value) {
return NumberUtils.isParsable(value.replace(',', '.'));
}
private static boolean isSimpleDouble(String valueAsString) {
return valueAsString.contains(".") && !valueAsString.contains("E") && !valueAsString.contains("e");
}
}

93
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java

@ -76,7 +76,7 @@ public class JsonConverter {
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement, long ts) throws JsonSyntaxException {
PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder();
convertToTelemetry(jsonElement, ts, null, builder, isTypeCastEnabled);
convertToTelemetry(jsonElement, ts, null, builder);
return builder.build();
}
@ -84,13 +84,13 @@ public class JsonConverter {
return convertToTelemetryProto(jsonElement, System.currentTimeMillis());
}
private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder, boolean typeCastEnabled) {
private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder) {
if (jsonElement.isJsonObject()) {
parseObject(systemTs, result, builder, jsonElement.getAsJsonObject(), typeCastEnabled);
parseObject(systemTs, result, builder, jsonElement.getAsJsonObject());
} else if (jsonElement.isJsonArray()) {
jsonElement.getAsJsonArray().forEach(je -> {
if (je.isJsonObject()) {
parseObject(systemTs, result, builder, je.getAsJsonObject(), typeCastEnabled);
parseObject(systemTs, result, builder, je.getAsJsonObject());
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + je);
}
@ -100,11 +100,11 @@ public class JsonConverter {
}
}
private static void parseObject(long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder, JsonObject jo, boolean typeCastEnabled) {
private static void parseObject(long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder, JsonObject jo) {
if (result != null) {
parseObject(result, systemTs, jo, typeCastEnabled);
parseObject(result, systemTs, jo);
} else {
parseObject(builder, systemTs, jo, typeCastEnabled);
parseObject(builder, systemTs, jo);
}
}
@ -146,7 +146,7 @@ public class JsonConverter {
public static PostAttributeMsg convertToAttributesProto(JsonElement jsonObject) throws JsonSyntaxException {
if (jsonObject.isJsonObject()) {
PostAttributeMsg.Builder result = PostAttributeMsg.newBuilder();
List<KeyValueProto> keyValueList = parseProtoValues(jsonObject.getAsJsonObject(), isTypeCastEnabled);
List<KeyValueProto> keyValueList = parseProtoValues(jsonObject.getAsJsonObject());
result.addAllKv(keyValueList);
return result.build();
} else {
@ -164,29 +164,29 @@ public class JsonConverter {
return result;
}
private static void parseObject(PostTelemetryMsg.Builder builder, long systemTs, JsonObject jo, boolean typeCastEnabled) {
private static void parseObject(PostTelemetryMsg.Builder builder, long systemTs, JsonObject jo) {
if (jo.has("ts") && jo.has("values")) {
parseWithTs(builder, jo, typeCastEnabled);
parseWithTs(builder, jo);
} else {
parseWithoutTs(builder, systemTs, jo, typeCastEnabled);
parseWithoutTs(builder, systemTs, jo);
}
}
private static void parseWithoutTs(PostTelemetryMsg.Builder request, long systemTs, JsonObject jo, boolean typeCastEnabled) {
private static void parseWithoutTs(PostTelemetryMsg.Builder request, long systemTs, JsonObject jo) {
TsKvListProto.Builder builder = TsKvListProto.newBuilder();
builder.setTs(systemTs);
builder.addAllKv(parseProtoValues(jo, typeCastEnabled));
builder.addAllKv(parseProtoValues(jo));
request.addTsKvList(builder.build());
}
private static void parseWithTs(PostTelemetryMsg.Builder request, JsonObject jo, boolean typeCastEnabled) {
private static void parseWithTs(PostTelemetryMsg.Builder request, JsonObject jo) {
TsKvListProto.Builder builder = TsKvListProto.newBuilder();
builder.setTs(jo.get("ts").getAsLong());
builder.addAllKv(parseProtoValues(jo.get("values").getAsJsonObject(), typeCastEnabled));
builder.addAllKv(parseProtoValues(jo.get("values").getAsJsonObject()));
request.addTsKvList(builder.build());
}
private static List<KeyValueProto> parseProtoValues(JsonObject valuesObject, boolean typeCastEnabled) {
private static List<KeyValueProto> parseProtoValues(JsonObject valuesObject) {
List<KeyValueProto> result = new ArrayList<>();
for (Entry<String, JsonElement> valueEntry : valuesObject.entrySet()) {
JsonElement element = valueEntry.getValue();
@ -197,9 +197,9 @@ public class JsonConverter {
String message = String.format("String value length [%d] for key [%s] is greater than maximum allowed [%d]", value.getAsString().length(), valueEntry.getKey(), maxStringValueLength);
throw new JsonSyntaxException(message);
}
if (typeCastEnabled && isNumber(value)) {
if (isTypeCastEnabled && NumberUtils.isParsable(value.getAsString())) {
try {
result.add(buildNumericKeyValueProto(value, valueEntry.getKey(), typeCastEnabled));
result.add(buildNumericKeyValueProto(value, valueEntry.getKey()));
} catch (RuntimeException th) {
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.STRING_V)
.setStringV(value.getAsString()).build());
@ -212,7 +212,7 @@ public class JsonConverter {
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.BOOLEAN_V)
.setBoolV(value.getAsBoolean()).build());
} else if (value.isNumber()) {
result.add(buildNumericKeyValueProto(value, valueEntry.getKey(), typeCastEnabled));
result.add(buildNumericKeyValueProto(value, valueEntry.getKey()));
} else if (!value.isJsonNull()) {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + value);
}
@ -225,15 +225,15 @@ public class JsonConverter {
return result;
}
private static KeyValueProto buildNumericKeyValueProto(JsonPrimitive value, String key, boolean typeCastEnabled) {
String valueAsString = value.getAsString().replace(',', '.');
private static KeyValueProto buildNumericKeyValueProto(JsonPrimitive value, String key) {
String valueAsString = value.getAsString();
KeyValueProto.Builder builder = KeyValueProto.newBuilder().setKey(key);
var bd = new BigDecimal(valueAsString);
if (bd.stripTrailingZeros().scale() <= 0 && !isSimpleDouble(valueAsString)) {
try {
return builder.setType(KeyValueType.LONG_V).setLongV(bd.longValueExact()).build();
} catch (ArithmeticException e) {
if (typeCastEnabled) {
if (isTypeCastEnabled) {
return builder.setType(KeyValueType.STRING_V).setStringV(bd.toPlainString()).build();
} else {
throw new JsonSyntaxException("Big integer values are not supported!");
@ -242,7 +242,7 @@ public class JsonConverter {
} else {
if (bd.scale() <= 16) {
return builder.setType(KeyValueType.DOUBLE_V).setDoubleV(bd.doubleValue()).build();
} else if (typeCastEnabled) {
} else if (isTypeCastEnabled) {
return builder.setType(KeyValueType.STRING_V).setStringV(bd.toPlainString()).build();
} else {
throw new JsonSyntaxException("Big integer values are not supported!");
@ -260,15 +260,15 @@ public class JsonConverter {
return TransportProtos.ToServerRpcRequestMsg.newBuilder().setRequestId(requestId).setMethodName(object.get("method").getAsString()).setParams(GSON.toJson(object.get("params"))).build();
}
private static void parseNumericValue(List<KvEntry> result, Entry<String, JsonElement> valueEntry, JsonPrimitive value, boolean typeCastEnabled) {
String valueAsString = value.getAsString().replace(',', '.');
private static void parseNumericValue(List<KvEntry> result, Entry<String, JsonElement> valueEntry, JsonPrimitive value) {
String valueAsString = value.getAsString();
String key = valueEntry.getKey();
var bd = new BigDecimal(valueAsString);
if (bd.stripTrailingZeros().scale() <= 0 && !isSimpleDouble(valueAsString)) {
try {
result.add(new LongDataEntry(key, bd.longValueExact()));
} catch (ArithmeticException e) {
if (typeCastEnabled) {
if (isTypeCastEnabled) {
result.add(new StringDataEntry(key, bd.toPlainString()));
} else {
throw new JsonSyntaxException("Big integer values are not supported!");
@ -277,7 +277,7 @@ public class JsonConverter {
} else {
if (bd.scale() <= 16) {
result.add(new DoubleDataEntry(key, bd.doubleValue()));
} else if (typeCastEnabled) {
} else if (isTypeCastEnabled) {
result.add(new StringDataEntry(key, bd.toPlainString()));
} else {
throw new JsonSyntaxException("Big integer values are not supported!");
@ -487,17 +487,13 @@ public class JsonConverter {
}
public static Set<AttributeKvEntry> convertToAttributes(JsonElement element) {
return convertToAttributes(element, isTypeCastEnabled);
}
public static Set<AttributeKvEntry> convertToAttributes(JsonElement element, boolean typeCastEnabled) {
Set<AttributeKvEntry> result = new HashSet<>();
long ts = System.currentTimeMillis();
result.addAll(parseValues(element.getAsJsonObject(), typeCastEnabled).stream().map(kv -> new BaseAttributeKvEntry(kv, ts)).collect(Collectors.toList()));
result.addAll(parseValues(element.getAsJsonObject()).stream().map(kv -> new BaseAttributeKvEntry(kv, ts)).collect(Collectors.toList()));
return result;
}
private static List<KvEntry> parseValues(JsonObject valuesObject, boolean typeCastEnabled) {
private static List<KvEntry> parseValues(JsonObject valuesObject) {
List<KvEntry> result = new ArrayList<>();
for (Entry<String, JsonElement> valueEntry : valuesObject.entrySet()) {
JsonElement element = valueEntry.getValue();
@ -508,9 +504,9 @@ public class JsonConverter {
String message = String.format("String value length [%d] for key [%s] is greater than maximum allowed [%d]", value.getAsString().length(), valueEntry.getKey(), maxStringValueLength);
throw new JsonSyntaxException(message);
}
if (typeCastEnabled && isNumber(value)) {
if (isTypeCastEnabled && NumberUtils.isParsable(value.getAsString())) {
try {
parseNumericValue(result, valueEntry, value, typeCastEnabled);
parseNumericValue(result, valueEntry, value);
} catch (RuntimeException th) {
result.add(new StringDataEntry(valueEntry.getKey(), value.getAsString()));
}
@ -520,7 +516,7 @@ public class JsonConverter {
} else if (value.isBoolean()) {
result.add(new BooleanDataEntry(valueEntry.getKey(), value.getAsBoolean()));
} else if (value.isNumber()) {
parseNumericValue(result, valueEntry, value, typeCastEnabled);
parseNumericValue(result, valueEntry, value);
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + value);
}
@ -545,35 +541,30 @@ public class JsonConverter {
public static Map<Long, List<KvEntry>> convertToTelemetry(JsonElement jsonElement, long systemTs, boolean sorted) throws
JsonSyntaxException {
return convertToTelemetry(jsonElement, systemTs, sorted, isTypeCastEnabled);
}
public static Map<Long, List<KvEntry>> convertToTelemetry(JsonElement jsonElement, long systemTs, boolean sorted, boolean typeCastEnabled) throws
JsonSyntaxException {
Map<Long, List<KvEntry>> result = sorted ? new TreeMap<>() : new HashMap<>();
convertToTelemetry(jsonElement, systemTs, result, null, typeCastEnabled);
convertToTelemetry(jsonElement, systemTs, result, null);
return result;
}
private static void parseObject(Map<Long, List<KvEntry>> result, long systemTs, JsonObject jo, boolean typeCastEnabled) {
private static void parseObject(Map<Long, List<KvEntry>> result, long systemTs, JsonObject jo) {
if (jo.has("ts") && jo.has("values")) {
parseWithTs(result, jo, typeCastEnabled);
parseWithTs(result, jo);
} else {
parseWithoutTs(result, systemTs, jo, typeCastEnabled);
parseWithoutTs(result, systemTs, jo);
}
}
private static void parseWithoutTs(Map<Long, List<KvEntry>> result, long systemTs, JsonObject jo, boolean typeCastEnabled) {
for (KvEntry entry : parseValues(jo, typeCastEnabled)) {
private static void parseWithoutTs(Map<Long, List<KvEntry>> result, long systemTs, JsonObject jo) {
for (KvEntry entry : parseValues(jo)) {
result.computeIfAbsent(systemTs, tmp -> new ArrayList<>()).add(entry);
}
}
public static void parseWithTs(Map<Long, List<KvEntry>> result, JsonObject jo, boolean typeCastEnabled) {
public static void parseWithTs(Map<Long, List<KvEntry>> result, JsonObject jo) {
long ts = jo.get("ts").getAsLong();
JsonObject valuesObject = jo.get("values").getAsJsonObject();
for (KvEntry entry : parseValues(valuesObject, typeCastEnabled)) {
for (KvEntry entry : parseValues(valuesObject)) {
result.computeIfAbsent(ts, tmp -> new ArrayList<>()).add(entry);
}
}
@ -643,10 +634,6 @@ public class JsonConverter {
}
private static boolean isNumber(JsonPrimitive value) {
return NumberUtils.isParsable(value.getAsString().replace(',', '.'));
}
private static String getStrValue(JsonObject jo, String field, boolean requiredField) {
if (jo.has(field)) {
return jo.get(field).getAsString();

6
common/transport/transport-api/src/test/java/JsonConverterTest.java

@ -65,12 +65,6 @@ public class JsonConverterTest {
Assert.assertEquals(1.1, result.get(0L).get(0).getDoubleValue().get(), 0.0);
}
@Test
public void testParseAsDoubleWithCommaDecimalSeparatorAndTypeCast() {
var result = JsonConverter.convertToTelemetry(JSON_PARSER.parse("{\"meterReadingDelta\": \"-1,1\"}"), 0L);
Assert.assertEquals(-1.1, result.get(0L).get(0).getDoubleValue().get(), 0.0);
}
@Test
public void testParseAsLong() {
var result = JsonConverter.convertToTelemetry(JSON_PARSER.parse("{\"meterReadingDelta\": 11}"), 0L);

Loading…
Cancel
Save