Browse Source

efento data point optimization: general parameters are saved once with the timestamp of the first sample of the first channel

pull/15333/head
dashevchenko 6 months ago
parent
commit
2cf1a76f03
  1. 17
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java
  2. 4
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/utils/CoapEfentoUtils.java
  3. 63
      common/transport/coap/src/test/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResourceTest.java

17
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java

@ -53,6 +53,7 @@ import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.TreeMap;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@ -258,6 +259,12 @@ public class CoapEfentoTransportResource extends AbstractCoapTransportResource {
}
Map<Long, JsonObject> valuesMap = new TreeMap<>();
// general measurements per message
Optional<ProtoChannel> channel1 = protoMeasurements.getChannelsList().stream()
.findFirst();
long startTs = channel1.map(value -> TimeUnit.SECONDS.toMillis(value.getTimestamp())).orElseGet(System::currentTimeMillis);
valuesMap.put(startTs, CoapEfentoUtils.setDefaultMeasurements(serialNumber, batteryStatus, nextTransmissionAtMillis, signal));
for (int channel = 0; channel < channelsList.size(); channel++) {
ProtoChannel protoChannel = channelsList.get(channel);
List<Integer> sampleOffsetsList = protoChannel.getSampleOffsetsList();
@ -271,6 +278,10 @@ public class CoapEfentoTransportResource extends AbstractCoapTransportResource {
long measurementPeriodMillis = TimeUnit.SECONDS.toMillis(measurementPeriod);
long startTimestampMillis = TimeUnit.SECONDS.toMillis(protoChannel.getTimestamp());
// measurements per channel
JsonObject tsValues = valuesMap.computeIfAbsent(startTimestampMillis, k -> new JsonObject());
tsValues.addProperty("measurement_interval", measurementPeriod);
for (int i = 0; i < sampleOffsetsList.size(); i++) {
int sampleOffset = sampleOffsetsList.get(i);
if (isSensorError(sampleOffset)) {
@ -290,13 +301,11 @@ public class CoapEfentoTransportResource extends AbstractCoapTransportResource {
}
long sampleOffsetMillis = TimeUnit.SECONDS.toMillis(sampleOffset);
long measurementTimestamp = startTimestampMillis + Math.abs(sampleOffsetMillis);
values = valuesMap.computeIfAbsent(measurementTimestamp - 1000, k ->
CoapEfentoUtils.setDefaultMeasurements(serialNumber, batteryStatus, measurementPeriod, nextTransmissionAtMillis, signal, k));
values = valuesMap.computeIfAbsent(measurementTimestamp - 1000, k -> new JsonObject());
addBinarySample(protoChannel, currentIsOk, values, channel + 1, sessionId);
} else {
long timestampMillis = startTimestampMillis + i * measurementPeriodMillis;
values = valuesMap.computeIfAbsent(timestampMillis, k -> CoapEfentoUtils.setDefaultMeasurements(
serialNumber, batteryStatus, measurementPeriod, nextTransmissionAtMillis, signal, k));
values = valuesMap.computeIfAbsent(timestampMillis, k -> new JsonObject());
addContinuesSample(protoChannel, sampleOffset, values, channel + 1, sessionId);
}
}

4
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/utils/CoapEfentoUtils.java

@ -59,14 +59,12 @@ public class CoapEfentoUtils {
return String.format("%s UTC", simpleDateFormat.format(new Date(timestampInMillis)));
}
public static JsonObject setDefaultMeasurements(String serialNumber, boolean batteryStatus, long measurementPeriod, long nextTransmissionAtMillis, long signal, long startTimestampMillis) {
public static JsonObject setDefaultMeasurements(String serialNumber, boolean batteryStatus, long nextTransmissionAtMillis, long signal) {
JsonObject values = new JsonObject();
values.addProperty("serial", serialNumber);
values.addProperty("battery", batteryStatus ? "ok" : "low");
values.addProperty("measured_at", convertTimestampToUtcString(startTimestampMillis));
values.addProperty("next_transmission_at", convertTimestampToUtcString(nextTransmissionAtMillis));
values.addProperty("signal", signal);
values.addProperty("measurement_interval", measurementPeriod);
return values;
}

63
common/transport/coap/src/test/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResourceTest.java

@ -36,10 +36,12 @@ import java.nio.ByteBuffer;
import java.text.SimpleDateFormat;
import java.time.Instant;
import java.util.Arrays;
import java.util.Comparator;
import java.util.Date;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
@ -121,7 +123,7 @@ class CoapEfentoTransportResourceTest {
assertThat(efentoMeasurements.get(1).getTs()).isEqualTo((tsInSec + 180 * 5) * 1000);
assertThat(efentoMeasurements.get(1).getValues().getAsJsonObject().get("temperature_1").getAsDouble()).isEqualTo(22.4);
assertThat(efentoMeasurements.get(1).getValues().getAsJsonObject().get("humidity_2").getAsDouble()).isEqualTo(30);
checkDefaultMeasurements(measurements, efentoMeasurements, 180 * 5, false);
checkDefaultMeasurements(measurements, efentoMeasurements, 180 * 5);
}
@ParameterizedTest
@ -149,7 +151,7 @@ class CoapEfentoTransportResourceTest {
assertThat(efentoMeasurements).hasSize(1);
assertThat(efentoMeasurements.get(0).getTs()).isEqualTo(tsInSec * 1000);
assertThat(efentoMeasurements.get(0).getValues().getAsJsonObject().get(property).getAsDouble()).isEqualTo(expectedValue);
checkDefaultMeasurements(measurements, efentoMeasurements, 180, false);
checkDefaultMeasurements(measurements, efentoMeasurements, 180);
}
private static Stream<Arguments> checkContinuousSensor() {
@ -207,7 +209,7 @@ class CoapEfentoTransportResourceTest {
assertThat(efentoMeasurements).hasSize(1);
assertThat(efentoMeasurements.get(0).getTs()).isEqualTo(tsInSec * 1000);
assertThat(efentoMeasurements.get(0).getValues().getAsJsonObject().get(totalPropertyName + "_2").getAsDouble()).isEqualTo(expectedTotalValue);
checkDefaultMeasurements(measurements, efentoMeasurements, 180, false);
checkDefaultMeasurements(measurements, efentoMeasurements, 180);
}
private static Stream<Arguments> checkPulseCounterSensors() {
@ -246,7 +248,7 @@ class CoapEfentoTransportResourceTest {
assertThat(efentoMeasurements).hasSize(1);
assertThat(efentoMeasurements.get(0).getTs()).isEqualTo(tsInSec * 1000);
assertThat(efentoMeasurements.get(0).getValues().getAsJsonObject().get("ok_alarm_1").getAsString()).isEqualTo("ALARM");
checkDefaultMeasurements(measurements, efentoMeasurements, 180 * 14, true);
checkDefaultMeasurements(measurements, efentoMeasurements, 180 * 14);
}
@ParameterizedTest
@ -275,7 +277,7 @@ class CoapEfentoTransportResourceTest {
assertThat(efentoMeasurements.get(0).getValues().getAsJsonObject().get(property).getAsString()).isEqualTo(expectedValueWhenOffsetNotOk);
assertThat(efentoMeasurements.get(1).getTs()).isEqualTo((tsInSec + 9) * 1000);
assertThat(efentoMeasurements.get(1).getValues().getAsJsonObject().get(property).getAsString()).isEqualTo(expectedValueWhenOffsetOk);
checkDefaultMeasurements(measurements, efentoMeasurements, 180, true);
checkDefaultMeasurements(measurements, efentoMeasurements, 180);
}
private static Stream<Arguments> checkBinarySensorWhenValueIsVarying() {
@ -306,31 +308,6 @@ class CoapEfentoTransportResourceTest {
.hasMessage("[" + sessionId + "]: Failed to get Efento measurements, reason: channels list is empty!");
}
@Test
void checkExceptionWhenValuesMapIsEmpty() {
long tsInSec = Instant.now().getEpochSecond();
ProtoMeasurements measurements = ProtoMeasurements.newBuilder()
.setSerialNumber(integerToByteString(1234))
.setCloudToken("test_token")
.setMeasurementPeriodBase(180)
.setMeasurementPeriodFactor(1)
.setBatteryStatus(true)
.setSignal(0)
.setNextTransmissionAt(1000)
.setTransferReason(0)
.setConfigurationHash(0)
.addChannels(MeasurementsProtos.ProtoChannel.newBuilder()
.setType(MEASUREMENT_TYPE_TEMPERATURE)
.setTimestamp(Math.toIntExact(tsInSec))
.build())
.build();
UUID sessionId = UUID.randomUUID();
assertThatThrownBy(() -> coapEfentoTransportResource.getEfentoMeasurements(measurements, sessionId))
.isInstanceOf(IllegalStateException.class)
.hasMessage("[" + sessionId + "]: Failed to collect Efento measurements, reason, values map is empty!");
}
// -------------------------------------------------------------------------
// ProtoDeviceInfo parsing tests
// -------------------------------------------------------------------------
@ -930,21 +907,17 @@ class CoapEfentoTransportResourceTest {
private void checkDefaultMeasurements(ProtoMeasurements incomingMeasurements,
List<CoapEfentoTransportResource.EfentoTelemetry> actualEfentoMeasurements,
long expectedMeasurementInterval,
boolean isBinarySensor) {
for (int i = 0; i < actualEfentoMeasurements.size(); i++) {
CoapEfentoTransportResource.EfentoTelemetry actualEfentoMeasurement = actualEfentoMeasurements.get(i);
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("serial").getAsString()).isEqualTo(CoapEfentoUtils.convertByteArrayToString(incomingMeasurements.getSerialNumber().toByteArray()));
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("battery").getAsString()).isEqualTo(incomingMeasurements.getBatteryStatus() ? "ok" : "low");
MeasurementsProtos.ProtoChannel protoChannel = incomingMeasurements.getChannelsList().get(0);
long measuredAt = isBinarySensor ?
TimeUnit.SECONDS.toMillis(protoChannel.getTimestamp()) + Math.abs(TimeUnit.SECONDS.toMillis(protoChannel.getSampleOffsetsList().get(i))) - 1000 :
TimeUnit.SECONDS.toMillis(protoChannel.getTimestamp() + i * expectedMeasurementInterval);
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("measured_at").getAsString()).isEqualTo(convertTimestampToUtcString(measuredAt));
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("next_transmission_at").getAsString()).isEqualTo(convertTimestampToUtcString(TimeUnit.SECONDS.toMillis(incomingMeasurements.getNextTransmissionAt())));
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("signal").getAsLong()).isEqualTo(incomingMeasurements.getSignal());
assertThat(actualEfentoMeasurement.getValues().getAsJsonObject().get("measurement_interval").getAsDouble()).isEqualTo(expectedMeasurementInterval);
}
long expectedMeasurementInterval) {
CoapEfentoTransportResource.EfentoTelemetry efentoTelemetry = actualEfentoMeasurements.stream()
.sorted(Comparator.comparing(CoapEfentoTransportResource.EfentoTelemetry::getTs))
.toList().get(0);
assertThat(efentoTelemetry.getValues().getAsJsonObject().get("serial").getAsString()).isEqualTo(CoapEfentoUtils.convertByteArrayToString(incomingMeasurements.getSerialNumber().toByteArray()));
assertThat(efentoTelemetry.getValues().getAsJsonObject().get("battery").getAsString()).isEqualTo(incomingMeasurements.getBatteryStatus() ? "ok" : "low");
assertThat(efentoTelemetry.getValues().getAsJsonObject().get("next_transmission_at").getAsString()).isEqualTo(convertTimestampToUtcString(TimeUnit.SECONDS.toMillis(incomingMeasurements.getNextTransmissionAt())));
assertThat(efentoTelemetry.getValues().getAsJsonObject().get("signal").getAsLong()).isEqualTo(incomingMeasurements.getSignal());
assertThat(efentoTelemetry.getValues().getAsJsonObject().get("measurement_interval").getAsDouble()).isEqualTo(expectedMeasurementInterval);
}
}

Loading…
Cancel
Save