Browse Source

sparkplug: add converter kvProto to value Metric array Bytes only

pull/8032/head
nickAS21 4 years ago
parent
commit
106fd69f5c
  1. 114
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java
  2. 41
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java
  3. 234
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java

114
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java

@ -41,28 +41,15 @@ import java.util.concurrent.atomic.AtomicReference;
import static org.awaitility.Awaitility.await;
import static org.thingsboard.common.util.JacksonUtil.newArrayNode;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.BooleanArray;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Bytes;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.DateTimeArray;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.DoubleArray;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.FloatArray;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int8;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int8Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.StringArray;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt16;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt16Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8Array;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric;
/**
@ -270,58 +257,6 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
String keys = "MyBytesArray";
byte[] bytes = {nextInt8(), nextInt8(), nextInt8()};
createdAddMetricTsKvJson(dataPayload, keys, bytes, ts, Bytes, listTsKvEntry, listKeys);
keys = "MyInt8Array";
Byte[] int8s = {nextInt8(), nextInt8(), nextInt8()};
createdAddMetricTsKvJson(dataPayload, keys, int8s, ts, Int8Array, listTsKvEntry, listKeys);
keys = "MyInt16Array";
Short[] int16s = {nextInt16(), nextInt16(), nextInt16()};
createdAddMetricTsKvJson(dataPayload, keys, int16s, ts, Int16Array, listTsKvEntry, listKeys);
keys = "MyUInt8Array";
Short[] uInt8s = {nextUInt8(), nextUInt8(), nextUInt8()};
createdAddMetricTsKvJson(dataPayload, keys, uInt8s, ts, UInt8Array, listTsKvEntry, listKeys);
keys = "MyInt32Array";
Integer[] int32s = {nextInt32(), nextInt32(), nextInt32()};
createdAddMetricTsKvJson(dataPayload, keys, int32s, ts, Int32Array, listTsKvEntry, listKeys);
keys = "MyUInt16Array";
Integer[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()};
createdAddMetricTsKvJson(dataPayload, keys, uInt16s, ts, UInt16Array, listTsKvEntry, listKeys);
keys = "MyInt64Array";
Long[] int64s = {nextInt64(), nextInt64(), nextInt64()};
createdAddMetricTsKvJson(dataPayload, keys, int64s, ts, Int64Array, listTsKvEntry, listKeys);
keys = "MyUInt32LArray";
Long[] uInt32Ls = {nextUInt32(), nextUInt32(), nextUInt32()};
createdAddMetricTsKvJson(dataPayload, keys, uInt32Ls, ts, UInt32Array, listTsKvEntry, listKeys);
keys = "MyDateTimeArray";
Long[] dateTimes = {nextDateTime(), nextDateTime(), nextDateTime()};
createdAddMetricTsKvJson(dataPayload, keys, dateTimes, ts, DateTimeArray, listTsKvEntry, listKeys);
keys = "MyUInt64Array";
Long[] uInt64s = {nextUInt64(), nextUInt64(), nextUInt64()};
createdAddMetricTsKvJson(dataPayload, keys, uInt64s, ts, UInt64Array, listTsKvEntry, listKeys);
keys = "MyFloatArray";
Float[] floats = {nextFloat(0, 300), nextFloat(0, 4000), nextFloat(10, 10000)};
createdAddMetricTsKvJson(dataPayload, keys, floats, ts, FloatArray, listTsKvEntry, listKeys);
keys = "MyDoubleArray";
Double[] doubles = {nextDouble(), nextDouble(), nextDouble()};
createdAddMetricTsKvJson(dataPayload, keys, doubles, ts, DoubleArray, listTsKvEntry, listKeys);
keys = "MyBooleanArray";
Boolean[] booleans = {nextBoolean(), nextBoolean(), nextBoolean()};
createdAddMetricTsKvJson(dataPayload, keys, booleans, ts, BooleanArray, listTsKvEntry, listKeys);
keys = "MyStringArray";
String[] strings = {nexString(), nexString(), nexString()};
createdAddMetricTsKvJson(dataPayload, keys, strings, ts, StringArray, listTsKvEntry, listKeys);
}
private TsKvEntry createdAddMetricTsKvLong(SparkplugBProto.Payload.Builder dataPayload, String key, Object value,
@ -372,54 +307,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
nodeArray.add(b);
}
break;
case Int8Array:
for (Byte b : (Byte[]) values) {
nodeArray.add(b);
}
break;
case Int16Array:
case UInt8Array:
for (Short b : (Short[]) values) {
nodeArray.add(b);
}
break;
case Int32Array:
case UInt16Array:
case Int64Array:
case UInt32Array:
case UInt64Array:
case DateTimeArray:
if (values instanceof Integer[]) {
for (Integer b : (Integer[]) values) {
nodeArray.add(b);
}
} else {
for (Long b : (Long[]) values) {
nodeArray.add(b);
}
}
break;
case DoubleArray:
for (Double b : (Double[]) values) {
nodeArray.add(b);
}
break;
case FloatArray:
for (Float b : (Float[]) values) {
nodeArray.add(b);
}
break;
case BooleanArray:
for (Boolean b : (Boolean[]) values) {
nodeArray.add(b);
}
break;
case StringArray:
for (String b : (String[]) values) {
nodeArray.add(b);
}
break;
default:
default:
throw new IllegalStateException("Unexpected value: " + metricDataType);
}
if (nodeArray.size() > 0) {

41
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java

@ -54,21 +54,6 @@ public enum MetricDataType {
// PropertyValue Types (20 and 21) are NOT metric datatypes
// Array Types
Int8Array(22, Byte[].class),
Int16Array(23, Short[].class),
Int32Array(24, Integer[].class),
Int64Array(25, Long[].class),
UInt8Array(26, Short[].class),
UInt16Array(27, Integer[].class),
UInt32Array(28, Long[].class),
UInt64Array(29, BigInteger[].class),
FloatArray(30, Float[].class),
DoubleArray(31, Double[].class),
BooleanArray(32, Boolean[].class),
StringArray(33, String[].class),
DateTimeArray(34, Date[].class),
// Unknown
Unknown(0, Object.class);
@ -155,32 +140,6 @@ public enum MetricDataType {
return File;
case 19:
return Template;
case 22:
return Int8Array;
case 23:
return Int16Array;
case 24:
return Int32Array;
case 25:
return Int64Array;
case 26:
return UInt8Array;
case 27:
return UInt16Array;
case 28:
return UInt32Array;
case 29:
return UInt64Array;
case 30:
return FloatArray;
case 31:
return DoubleArray;
case 32:
return BooleanArray;
case 33:
return StringArray;
case 34:
return DateTimeArray;
default:
return Unknown;
}

234
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java

@ -30,19 +30,8 @@ import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.DataOutputStream;
import java.io.IOException;
import java.io.ObjectInputStream;
import java.io.ObjectOutputStream;
import java.math.BigDecimal;
import java.nio.ByteBuffer;
import java.nio.DoubleBuffer;
import java.nio.FloatBuffer;
import java.nio.IntBuffer;
import java.nio.LongBuffer;
import java.nio.ShortBuffer;
import java.text.NumberFormat;
import java.util.Arrays;
import java.util.List;
@ -107,85 +96,13 @@ public class SparkplugMetricUtil {
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.STRING_V)
.setStringV(protoMetric.getStringValue()).build());
// byte[]
case BooleanArray:
ByteBuffer booleanByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray());
while (booleanByteBuffer.hasRemaining()) {
nodeArray.add(booleanByteBuffer.get() == (byte) 0 ? false : true);
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// byte[]
case Bytes:
case Int8Array:
ByteBuffer byteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray());
while (byteBuffer.hasRemaining()) {
nodeArray.add(byteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// short[]
case Int16Array:
case UInt8Array:
ShortBuffer shortByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray())
.asShortBuffer();
while (shortByteBuffer.hasRemaining()) {
nodeArray.add(shortByteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// int[]
case Int32Array:
case UInt16Array:
IntBuffer intByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray())
.asIntBuffer();
while (intByteBuffer.hasRemaining()) {
nodeArray.add(intByteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// float[]
case FloatArray:
FloatBuffer floatByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray())
.asFloatBuffer();
while (floatByteBuffer.hasRemaining()) {
nodeArray.add(floatByteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// double[]
case DoubleArray:
DoubleBuffer doubleByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray())
.asDoubleBuffer();
while (doubleByteBuffer.hasRemaining()) {
nodeArray.add(doubleByteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
// long[]
case DateTimeArray:
case Int64Array:
case UInt64Array:
case UInt32Array:
LongBuffer longByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray())
.asLongBuffer();
while (longByteBuffer.hasRemaining()) {
nodeArray.add(longByteBuffer.get());
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
case StringArray:
ByteBuffer stringByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray());
final ByteArrayInputStream byteArrayInputStream =
new ByteArrayInputStream(stringByteBuffer.array());
final ObjectInputStream objectInputStream =
new ObjectInputStream(byteArrayInputStream);
final String[] stringArray = (String[]) objectInputStream.readObject();
objectInputStream.close();
for (String s : stringArray) {
nodeArray.add(s);
}
return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(nodeArray.toString()).build());
case DataSet:
case Template:
case File:
@ -258,46 +175,6 @@ public class SparkplugMetricUtil {
case Bytes:
ByteString byteString = ByteString.copyFrom((byte[]) value);
return metric.toBuilder().setBytesValue(byteString).build();
case Int8Array:
byte [] int8Array = byteArrayToByteArray((Byte[]) value);
ByteString byteInt8Array = ByteString.copyFrom((int8Array));
return metric.toBuilder().setBytesValue(byteInt8Array).build();
case Int16Array:
case UInt8Array:
byte[] int16Array = shortArrayToByteArray((Short[]) value);
ByteString byteInt16Array = ByteString.copyFrom((int16Array));
return metric.toBuilder().setBytesValue(byteInt16Array).build();
case Int32Array:
case UInt16Array:
case Int64Array:
case UInt32Array:
case UInt64Array:
case DateTimeArray:
if (value instanceof Integer[]) {
byte[] int32Array = integerArrayToByteArray((Integer[]) value);
ByteString byteInt32Array = ByteString.copyFrom((int32Array));
return metric.toBuilder().setBytesValue(byteInt32Array).build();
} else {
byte[] int64Array = longArrayToByteArray((Long[]) value);
ByteString byteInt64Array = ByteString.copyFrom((int64Array));
return metric.toBuilder().setBytesValue(byteInt64Array).build();
}
case DoubleArray:
byte[] doubleArray = doublArrayToByteArray((Double[]) value);
ByteString byteDoubleArray = ByteString.copyFrom(doubleArray);
return metric.toBuilder().setBytesValue(byteDoubleArray).build();
case FloatArray:
byte[] floatArray = floatArrayToByteArray((Float[]) value);
ByteString byteFloatArray = ByteString.copyFrom(floatArray);
return metric.toBuilder().setBytesValue(byteFloatArray).build();
case BooleanArray:
byte[] booleanArray = booleanArrayToByteArray((Boolean[]) value);
ByteString byteBooleanArray = ByteString.copyFrom(booleanArray);
return metric.toBuilder().setBytesValue(byteBooleanArray).build();
case StringArray:
byte[] stringArray = stringArrayToByteArray((String[]) value);
ByteString byteStringArray = ByteString.copyFrom(stringArray);
return metric.toBuilder().setBytesValue(byteStringArray).build();
case DataSet:
return metric.toBuilder().setDatasetValue((SparkplugBProto.Payload.DataSet) value).build();
case File:
@ -417,43 +294,6 @@ public class SparkplugMetricUtil {
bytes[i] = listBytes.get(i).byteValue();
}
return Optional.of(bytes);
// Byte []
case Int8Array:
List<Byte> listInt8Array = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
Byte[] int8Arrays = listInt8Array.toArray(Byte[]::new) ;
return Optional.of(int8Arrays);
// Short[]
case Int16Array:
case UInt8Array:
List<Short> listShorts = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listShorts.toArray(Short[]::new));
// Integer []
case UInt16Array:
case Int32Array:
List<Integer> listIntegers = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listIntegers.toArray(Integer[]::new));
// Long[]
case UInt32Array:
case Int64Array:
case UInt64Array:
case DateTimeArray:
List<Long> listLongs = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listLongs.toArray(Long[]::new));
// Double []
case DoubleArray:
List<Double> listDoubles = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listDoubles.toArray(Double[]::new));
// Float[]
case FloatArray:
List<Float> listFloats = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listFloats.toArray(Float[]::new));
// Boolean[]
case BooleanArray:
List<Boolean> listBooleans = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listBooleans.toArray(Boolean[]::new));
case StringArray:
List<String> listStrings = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
return Optional.of(listStrings.toArray(String[]::new));
case DataSet:
case File:
case Template:
@ -470,80 +310,6 @@ public class SparkplugMetricUtil {
}
}
private static byte[] byteArrayToByteArray(Byte[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length);
for (Byte d : inputs) {
bb.put(d);
}
return bb.array();
}
private static byte[] shortArrayToByteArray(Short[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 2);
for (Short d : inputs) {
bb.putShort(d);
}
return bb.array();
}
private static byte[] integerArrayToByteArray(Integer[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4);
for (Integer d : inputs) {
bb.putInt(d);
}
return bb.array();
}
private static byte[] longArrayToByteArray(Long[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8);
for (Long d : inputs) {
bb.putLong(d);
}
return bb.array();
}
private static byte[] doublArrayToByteArray(Double[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8);
for (Double d : inputs) {
bb.putDouble(d);
}
return bb.array();
}
private static byte[] floatArrayToByteArray(Float[] inputs) throws ThingsboardException {
ByteArrayOutputStream bas = new ByteArrayOutputStream();
DataOutputStream ds = new DataOutputStream(bas);
for (Float f : inputs) {
try {
ds.writeFloat(f);
} catch (IOException e) {
throw new ThingsboardException("Invalid value float ", ThingsboardErrorCode.INVALID_ARGUMENTS);
}
}
return bas.toByteArray();
}
private static byte[] booleanArrayToByteArray(Boolean[] inputs) {
byte[] toReturn = new byte[inputs.length];
for (int entry = 0; entry < toReturn.length; entry++) {
toReturn[entry] = (byte) (inputs[entry] ? 1 : 0);
}
return toReturn;
}
private static byte[] stringArrayToByteArray(String[] inputs) throws ThingsboardException {
final ByteArrayOutputStream bas = new ByteArrayOutputStream();
try {
final ObjectOutputStream os = new ObjectOutputStream(bas);
os.writeObject(inputs);
os.flush();
os.close();
} catch (Exception e) {
throw new ThingsboardException("Invalid value float ", ThingsboardErrorCode.INVALID_ARGUMENTS);
}
return bas.toByteArray();
}
private static Optional<String> getValueKvProtoPrimitive(TransportProtos.KeyValueProto kv) {
if (kv.getTypeValue() == 0) { // boolean
return Optional.of(String.valueOf(kv.getBoolV()));

Loading…
Cancel
Save