From 38f6e314b1000c8c7d8a1b58fd3a2f2aa8ba1107 Mon Sep 17 00:00:00 2001 From: Viacheslav Klimov Date: Fri, 13 Aug 2021 13:23:56 +0300 Subject: [PATCH] Refactor --- .../service/asset/AssetBulkImportService.java | 11 ++- .../device/DeviceBulkImportService.java | 33 ++++--- .../service/edge/EdgeBulkImportService.java | 11 ++- .../importing/AbstractBulkImportService.java | 80 ++++++++++++---- .../importing/BulkImportColumnType.java | 18 +++- .../server/utils/TypeCastUtil.java | 55 +++++++++++ .../transport/adaptor/JsonConverter.java | 93 ++++++++----------- .../src/test/java/JsonConverterTest.java | 6 -- 8 files changed, 199 insertions(+), 108 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java diff --git a/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java index 54d5f942e2..2f94ac3b0d 100644 --- a/application/src/main/java/org/thingsboard/server/service/asset/AssetBulkImportService.java +++ b/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 { } @Override - protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map entityData, SecurityUser user) { + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { ImportedEntityInfo 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 { return importedEntityInfo; } - private void setAssetFields(Asset asset, Map data) { + private void setAssetFields(Asset asset, Map 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; diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java index 9799fb1a97..5f5e363b45 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java +++ b/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 { } @Override - protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map entityData, SecurityUser user) { + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { ImportedEntityInfo 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 = 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 { return importedEntityInfo; } - private void setDeviceFields(Device device, Map data) { + private void setDeviceFields(Device device, Map 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 { } @SneakyThrows - private DeviceCredentials createDeviceCredentials(Map data) { - Set columns = data.keySet().stream().map(BulkImportRequest.ColumnMapping::getType).collect(Collectors.toSet()); + private DeviceCredentials createDeviceCredentials(Map fields) { + Set 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 { 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 { 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 { 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)); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java index b1d85e5660..ec6a2a4e55 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeBulkImportService.java +++ b/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 { } @Override - protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map entityData, SecurityUser user) { + protected ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map fields, SecurityUser user) { ImportedEntityInfo 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 { return importedEntityInfo; } - private void setEdgeFields(Edge edge, Map data) { + private void setEdgeFields(Edge edge, Map 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; diff --git a/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java index 51798419ef..d042bc51b5 100644 --- a/application/src/main/java/org/thingsboard/server/service/importing/AbstractBulkImportService.java +++ b/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 { i.incrementAndGet(); try { - ImportedEntityInfo importedEntityInfo = saveEntity(request, entityData, user); + ImportedEntityInfo 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 saveEntity(BulkImportRequest importRequest, Map entityData, SecurityUser user); + protected abstract ImportedEntityInfo saveEntity(BulkImportRequest importRequest, Map 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 data) { - Stream.of(BulkImportColumnType.SHARED_ATTRIBUTE, BulkImportColumnType.SERVER_ATTRIBUTE, BulkImportColumnType.TIMESERIES) + private void saveKvs(SecurityUser user, E entity, Map 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 kvsEntry) { - List timeseries = JsonConverter.convertToTelemetry(kvsEntry.getValue(), System.currentTimeMillis(), false, true) + List 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 kvsEntry, BulkImportColumnType kvType) { String scope = kvType.getKey(); - List attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue(), true)); + List 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 data) { - return data.entrySet().stream().filter(entry -> entry.getKey().getType() == bulkImportColumnType).findFirst().map(Map.Entry::getValue).orElse(null); - } - - private List> parseData(BulkImportRequest request) throws Exception { + private List parseData(BulkImportRequest request) throws Exception { List> records = CsvUtils.parseCsv(request.getFile(), request.getMapping().getDelimiter()); if (request.getMapping().getHeader()) { records.remove(0); } List 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 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 fields = new HashMap<>(); + private final Map 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(); + } + + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java b/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java index a0ef0a15bc..f0f870b93c 100644 --- a/application/src/main/java/org/thingsboard/server/service/importing/BulkImportColumnType.java +++ b/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; + } } diff --git a/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java b/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java new file mode 100644 index 0000000000..f2b7219ce0 --- /dev/null +++ b/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 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"); + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java index 1db75193bf..be4143e388 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java @@ -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> result, PostTelemetryMsg.Builder builder, boolean typeCastEnabled) { + private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map> 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> result, PostTelemetryMsg.Builder builder, JsonObject jo, boolean typeCastEnabled) { + private static void parseObject(long systemTs, Map> 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 keyValueList = parseProtoValues(jsonObject.getAsJsonObject(), isTypeCastEnabled); + List 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 parseProtoValues(JsonObject valuesObject, boolean typeCastEnabled) { + private static List parseProtoValues(JsonObject valuesObject) { List result = new ArrayList<>(); for (Entry 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 result, Entry valueEntry, JsonPrimitive value, boolean typeCastEnabled) { - String valueAsString = value.getAsString().replace(',', '.'); + private static void parseNumericValue(List result, Entry 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 convertToAttributes(JsonElement element) { - return convertToAttributes(element, isTypeCastEnabled); - } - - public static Set convertToAttributes(JsonElement element, boolean typeCastEnabled) { Set 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 parseValues(JsonObject valuesObject, boolean typeCastEnabled) { + private static List parseValues(JsonObject valuesObject) { List result = new ArrayList<>(); for (Entry 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> convertToTelemetry(JsonElement jsonElement, long systemTs, boolean sorted) throws JsonSyntaxException { - return convertToTelemetry(jsonElement, systemTs, sorted, isTypeCastEnabled); - } - - public static Map> convertToTelemetry(JsonElement jsonElement, long systemTs, boolean sorted, boolean typeCastEnabled) throws - JsonSyntaxException { Map> result = sorted ? new TreeMap<>() : new HashMap<>(); - convertToTelemetry(jsonElement, systemTs, result, null, typeCastEnabled); + convertToTelemetry(jsonElement, systemTs, result, null); return result; } - private static void parseObject(Map> result, long systemTs, JsonObject jo, boolean typeCastEnabled) { + private static void parseObject(Map> 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> result, long systemTs, JsonObject jo, boolean typeCastEnabled) { - for (KvEntry entry : parseValues(jo, typeCastEnabled)) { + private static void parseWithoutTs(Map> result, long systemTs, JsonObject jo) { + for (KvEntry entry : parseValues(jo)) { result.computeIfAbsent(systemTs, tmp -> new ArrayList<>()).add(entry); } } - public static void parseWithTs(Map> result, JsonObject jo, boolean typeCastEnabled) { + public static void parseWithTs(Map> 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(); diff --git a/common/transport/transport-api/src/test/java/JsonConverterTest.java b/common/transport/transport-api/src/test/java/JsonConverterTest.java index 8901afe5c2..2c1a3c5551 100644 --- a/common/transport/transport-api/src/test/java/JsonConverterTest.java +++ b/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);