Browse Source

convert ValueWithTs to a record & replace usages of Collections.singletonList to List.of & use single key search for async method

pull/10483/head
ShvaykaD 3 years ago
parent
commit
85d229be2d
  1. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  2. 34
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java
  3. 11
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java

2
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -38,7 +38,7 @@ public interface TimeseriesService {
ListenableFuture<List<TsKvEntry>> findAll(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries);
ListenableFuture<Optional<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, String keys);
ListenableFuture<Optional<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, String key);
ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys);

34
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java

@ -37,7 +37,6 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@ -129,8 +128,8 @@ public class CalculateDeltaNode implements TbNode {
}
private ListenableFuture<ValueWithTs> fetchLatestValueAsync(EntityId entityId) {
return Futures.transform(timeseriesService.findLatest(ctx.getTenantId(), entityId, Collections.singletonList(config.getInputValueKey())),
list -> extractValue(list.get(0))
return Futures.transform(timeseriesService.findLatest(ctx.getTenantId(), entityId, config.getInputValueKey()),
tsKvEntryOpt -> tsKvEntryOpt.map(this::extractValue).orElse(null)
, ctx.getDbCallbackExecutor());
}
@ -138,7 +137,7 @@ public class CalculateDeltaNode implements TbNode {
List<TsKvEntry> tsKvEntries = timeseriesService.findLatestSync(
ctx.getTenantId(),
entityId,
Collections.singletonList(config.getInputValueKey()));
List.of(config.getInputValueKey()));
return extractValue(tsKvEntries.get(0));
}
@ -161,36 +160,23 @@ public class CalculateDeltaNode implements TbNode {
double result = 0.0;
long ts = kvEntry.getTs();
switch (kvEntry.getDataType()) {
case LONG:
result = kvEntry.getLongValue().get();
break;
case DOUBLE:
result = kvEntry.getDoubleValue().get();
break;
case STRING:
case LONG -> result = kvEntry.getLongValue().get();
case DOUBLE -> result = kvEntry.getDoubleValue().get();
case STRING -> {
try {
result = Double.parseDouble(kvEntry.getStrValue().get());
} catch (NumberFormatException e) {
throw new IllegalArgumentException("Calculation failed. Unable to parse value [" + kvEntry.getStrValue().get() + "]" +
" of telemetry [" + kvEntry.getKey() + "] to Double");
}
break;
case BOOLEAN:
throw new IllegalArgumentException("Calculation failed. Boolean values are not supported!");
case JSON:
throw new IllegalArgumentException("Calculation failed. JSON values are not supported!");
}
case BOOLEAN -> throw new IllegalArgumentException("Calculation failed. Boolean values are not supported!");
case JSON -> throw new IllegalArgumentException("Calculation failed. JSON values are not supported!");
}
return new ValueWithTs(ts, result);
}
private static class ValueWithTs {
private final long ts;
private final double value;
private ValueWithTs(long ts, double value) {
this.ts = ts;
this.value = value;
}
private record ValueWithTs(long ts, double value) {
}
}

11
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java

@ -47,6 +47,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import java.util.List;
import java.util.Optional;
import java.util.UUID;
import static org.junit.jupiter.api.Assertions.assertEquals;
@ -68,9 +69,9 @@ import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class CalculateDeltaNodeTest {
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID());
private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID());
private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
private final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.fromString("2ba3ded4-882b-40cf-999a-89da9ccd58f9"));
private final TenantId TENANT_ID = new TenantId(UUID.fromString("3842e740-0d89-43a9-8d52-ae44023847ba"));
private final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
@Mock
private TbContext ctxMock;
@Mock
@ -435,8 +436,8 @@ public class CalculateDeltaNodeTest {
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(timeseriesServiceMock.findLatest(
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey())))
)).thenReturn(Futures.immediateFuture(List.of(tsKvEntry)));
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), eq(tsKvEntry.getKey())
)).thenReturn(Futures.immediateFuture(Optional.of(tsKvEntry)));
}
@RequiredArgsConstructor

Loading…
Cancel
Save