Browse Source

Merge pull request #14549 from volodymyr-babak/edge-stats-test

Improve Edge stats reporting and add related integration tests
pull/14612/head
Viacheslav Klimov 10 months ago
committed by GitHub
parent
commit
ba5b90db6d
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 1
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 12
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java
  4. 13
      application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java
  5. 148
      application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java
  6. 133
      application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java
  7. 1
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java
  8. 3
      dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java
  9. 15
      dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java

1
application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java

@ -206,7 +206,6 @@ public class EdgeContextComponent {
@Autowired @Autowired
private Optional<EdgeStatsCounterService> statsCounterService; private Optional<EdgeStatsCounterService> statsCounterService;
// processors // processors
@Autowired @Autowired
private AlarmProcessor alarmProcessor; private AlarmProcessor alarmProcessor;

12
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -306,7 +306,8 @@ public abstract class EdgeGrpcSession implements Closeable {
if (isConnected() && !pageData.getData().isEmpty()) { if (isConnected() && !pageData.getData().isEmpty()) {
if (fetcher instanceof GeneralEdgeEventFetcher) { if (fetcher instanceof GeneralEdgeEventFetcher) {
long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage()); long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage());
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.setDownlinkMsgsLag(edge.getTenantId(), edge.getId(), queueSize)); ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, edge.getId(), queueSize));
} }
log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size()); log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size());
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData());
@ -504,7 +505,8 @@ public abstract class EdgeGrpcSession implements Closeable {
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId())
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg)
.error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); .error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build());
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), edge.getId(), copy.size())); ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), edge.getId(), copy.size()));
stopCurrentSendDownlinkMsgsTask(false); stopCurrentSendDownlinkMsgsTask(false);
} }
} else { } else {
@ -544,7 +546,8 @@ public abstract class EdgeGrpcSession implements Closeable {
try { try {
if (msg.getSuccess()) { if (msg.getSuccess()) {
sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId());
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PUSHED, edge.getTenantId(), edge.getId(), 1)); ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PUSHED, edge.getTenantId(), edge.getId(), 1));
log.debug("[{}][{}][{}] Msg has been processed successfully! Msg Id: [{}], Msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg); log.debug("[{}][{}][{}] Msg has been processed successfully! Msg Id: [{}], Msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg);
} else { } else {
log.debug("[{}][{}][{}] Msg processing failed! Msg Id: [{}], Error msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg.getErrorMsg()); log.debug("[{}][{}][{}] Msg processing failed! Msg Id: [{}], Error msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg.getErrorMsg());
@ -813,7 +816,8 @@ public abstract class EdgeGrpcSession implements Closeable {
} }
} }
highPriorityQueue.add(edgeEvent); highPriorityQueue.add(edgeEvent);
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edge.getTenantId(), edgeEvent.getEdgeId(), 1)); ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edge.getTenantId(), edgeEvent.getEdgeId(), 1));
} }
protected ListenableFuture<List<Void>> processUplinkMsg(UplinkMsg uplinkMsg) { protected ListenableFuture<List<Void>> processUplinkMsg(UplinkMsg uplinkMsg) {

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java

@ -52,7 +52,8 @@ public class KafkaEdgeEventService extends BaseEdgeEventService {
TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId()); TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId());
ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build(); ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build();
producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null);
statsCounterService.ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }

13
application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java

@ -54,7 +54,7 @@ import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_P
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_TMP_FAILED; import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_TMP_FAILED;
@TbCoreComponent @TbCoreComponent
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true", matchIfMissing = false) @ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true")
@RequiredArgsConstructor @RequiredArgsConstructor
@Service @Service
@Slf4j @Slf4j
@ -70,7 +70,6 @@ public class EdgeStatsService {
@Value("${edges.stats.report-interval-millis:600000}") @Value("${edges.stats.report-interval-millis:600000}")
private long reportIntervalMillis; private long reportIntervalMillis;
@Scheduled( @Scheduled(
fixedDelayString = "${edges.stats.report-interval-millis:600000}", fixedDelayString = "${edges.stats.report-interval-millis:600000}",
initialDelayString = "${edges.stats.report-interval-millis:600000}" initialDelayString = "${edges.stats.report-interval-millis:600000}"
@ -80,13 +79,13 @@ public class EdgeStatsService {
long now = System.currentTimeMillis(); long now = System.currentTimeMillis();
long ts = now - (now % reportIntervalMillis); long ts = now - (now % reportIntervalMillis);
Map<EdgeId, MsgCounters> countersByEdge = statsCounterService.getCounterByEdge(); Map<EdgeId, MsgCounters> countersByEdgeSnapshot = new HashMap<>(statsCounterService.getMsgCountersByEdge());
Map<EdgeId, Long> lagByEdgeId = kafkaAdmin.isPresent() ? getEdgeLagByEdgeId(countersByEdge) : Collections.emptyMap(); boolean isKafkaStats = kafkaAdmin.isPresent();
Map<EdgeId, MsgCounters> countersByEdgeSnapshot = new HashMap<>(statsCounterService.getCounterByEdge()); Map<EdgeId, Long> lagByEdgeId = isKafkaStats ? getLagByEdgeId(countersByEdgeSnapshot) : Collections.emptyMap();
countersByEdgeSnapshot.forEach((edgeId, counters) -> { countersByEdgeSnapshot.forEach((edgeId, counters) -> {
TenantId tenantId = counters.getTenantId(); TenantId tenantId = counters.getTenantId();
if (kafkaAdmin.isPresent()) { if (isKafkaStats) {
counters.getMsgsLag().set(lagByEdgeId.getOrDefault(edgeId, 0L)); counters.getMsgsLag().set(lagByEdgeId.getOrDefault(edgeId, 0L));
} }
List<TsKvEntry> statsEntries = List.of( List<TsKvEntry> statsEntries = List.of(
@ -102,7 +101,7 @@ public class EdgeStatsService {
}); });
} }
private Map<EdgeId, Long> getEdgeLagByEdgeId(Map<EdgeId, MsgCounters> countersByEdge) { private Map<EdgeId, Long> getLagByEdgeId(Map<EdgeId, MsgCounters> countersByEdge) {
Map<EdgeId, String> edgeToTopicMap = countersByEdge.entrySet().stream() Map<EdgeId, String> edgeToTopicMap = countersByEdge.entrySet().stream()
.collect(Collectors.toMap( .collect(Collectors.toMap(
Map.Entry::getKey, Map.Entry::getKey,

148
application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java

@ -0,0 +1,148 @@
/**
* Copyright © 2016-2025 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.edge;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.extern.slf4j.Slf4j;
import org.junit.Assert;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService;
import org.thingsboard.server.dao.edge.stats.EdgeStatsKey;
import org.thingsboard.server.dao.edge.stats.MsgCounters;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.service.edge.stats.EdgeStatsService;
import java.time.Duration;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import static org.awaitility.Awaitility.await;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_ADDED;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_PUSHED;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_TMP_FAILED;
@DaoSqlTest
@Slf4j
public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
private static final String STATISTICS_DEVICE_PROFILE = "STATISTICS";
private static final long EXPECTED_MSGS_ADDED = 6L;
private static final long EXPECTED_MSGS_PUSHED = 6L;
private static final long EXPECTED_MSGS_PERMANENTLY_FAILED = 0L;
private static final long EXPECTED_MSGS_TMP_FAILED = 0L;
@Autowired
private EdgeStatsService edgeStatsService;
@Autowired
private EdgeStatsCounterService statsCounterService;
@Test
public void testReportStats() throws Exception {
// GIVEN
simulateEdgeEventsAddedDownlinkPushed();
// Await Edge Counters Updated
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
MsgCounters counters = statsCounterService.getMsgCountersByEdge().get(edge.getId());
assertEquals(EXPECTED_MSGS_ADDED, counters.getMsgsAdded().get());
assertEquals(EXPECTED_MSGS_PUSHED, counters.getMsgsPushed().get());
assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, counters.getMsgsPermanentlyFailed().get());
assertEquals(EXPECTED_MSGS_TMP_FAILED, counters.getMsgsTmpFailed().get());
});
Thread.sleep(1000);
// WHEN
edgeStatsService.reportStats();
// THEN
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
List<TsKvEntry> actualStats = fetchLatestStats();
assertEquals(EXPECTED_MSGS_ADDED, getStatsLongValue(actualStats, DOWNLINK_MSGS_ADDED));
assertEquals(EXPECTED_MSGS_PUSHED, getStatsLongValue(actualStats, DOWNLINK_MSGS_PUSHED));
assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, getStatsLongValue(actualStats, DOWNLINK_MSGS_PERMANENTLY_FAILED));
assertEquals(EXPECTED_MSGS_TMP_FAILED, getStatsLongValue(actualStats, DOWNLINK_MSGS_TMP_FAILED));
});
}
private long getStatsLongValue(List<TsKvEntry> stats, EdgeStatsKey key) {
return stats.stream().filter(e -> e.getKey().equals(key.getKey())).findFirst().get().getLongValue().orElse(0L);
}
private List<TsKvEntry> fetchLatestStats() throws ExecutionException, InterruptedException {
return tsService.findLatest(
tenantId,
edge.getId(),
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()).get();
}
private void simulateEdgeEventsAddedDownlinkPushed() throws InterruptedException, ExecutionException {
statsCounterService.clear(edge.getId());
// Save device and assign to edge
// 2 DOWNLINK_MSGS_ADDED, EdgeEvents: [{DEVICE_PROFILE: ADDED}, {DEVICE: ASSIGNED_TO_EDGE}]
// 2 DOWNLINK_MSGS_PUSHED, Downlinks: [{deviceProfileUpdateMsg}, {deviceUpdateMsg, deviceProfileUpdateMsg, deviceCredentialsUpdateMsg}]
edgeImitator.expectMessageAmount(4);
Device savedDevice = saveDevice("StatisticDevice", STATISTICS_DEVICE_PROFILE);
doPost("/api/edge/" + edge.getUuidId() + "/device/" + savedDevice.getUuidId(), Device.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// Save asset and assign to edge
// 1 DOWNLINK_MSGS_ADDED, EdgeEvents: [{ASSET: ASSIGNED_TO_EDGE}]
// 1 DOWNLINK_MSGS_PUSHED, Downlinks: [{assetUpdateMsg, assetProfileUpdateMsg}]
edgeImitator.expectMessageAmount(2);
Asset savedAsset = saveAsset("Edge Asset");
doPost("/api/edge/" + edge.getUuidId()
+ "/asset/" + savedAsset.getUuidId(), Asset.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// Create customer and assign edge to the customer
// 2 DOWNLINK_MSGS_ADDED, EdgeEvents: [{CUSTOMER: ADDED}, {EDGE: ASSIGNED_TO_CUSTOMER}]
// 2 DOWNLINK_MSGS_PUSHED, Downlinks: [{customerUpdateMsg}, {edgeConfiguration}]
edgeImitator.expectMessageAmount(2);
Customer customer = new Customer();
customer.setTitle("Edge Customer");
Customer savedCustomer = doPost("/api/customer", customer, Customer.class);
doPost("/api/customer/" + savedCustomer.getUuidId()
+ "/edge/" + edge.getUuidId(), Edge.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// Send device telemetry downlink for the device
// 1 DOWNLINK_MSGS_ADDED, EdgeEvents: [{DEVICE: TIMESERIES_UPDATED}]
// 1 DOWNLINK_MSGS_PUSHED, Downlinks: [{entityData}]
edgeImitator.expectMessageAmount(1);
String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = JacksonUtil.toJsonNode(timeseriesData);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, savedDevice.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeEventService.saveAsync(edgeEvent).get();
Assert.assertTrue(edgeImitator.waitForMessages());
}
}

133
application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java

@ -21,6 +21,7 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils; import org.springframework.test.util.ReflectionTestUtils;
@ -44,9 +45,11 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when; import static org.mockito.Mockito.when;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_ADDED; import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_ADDED;
@ -58,6 +61,16 @@ import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_T
@ExtendWith(MockitoExtension.class) @ExtendWith(MockitoExtension.class)
public class EdgeStatsTest { public class EdgeStatsTest {
private static final int TTL_DAYS = 30;
private static final long REPORT_INTERVAL_MILLIS = 600_000L;
private static final long EXPECTED_MSGS_ADDED = 5L;
private static final long EXPECTED_MSGS_PUSHED = 3L;
private static final long EXPECTED_MSGS_PERMANENTLY_FAILED = 1L;
private static final long EXPECTED_MSGS_TMP_FAILED = 0L;
private static final long EXPECTED_MSGS_LAG = 10L;
private static final long EXPECTED_MSGS_KAFKA_LAG = 15L;
@Mock @Mock
private TimeseriesService tsService; private TimeseriesService tsService;
@Mock @Mock
@ -66,108 +79,100 @@ public class EdgeStatsTest {
private EdgeStatsCounterService statsCounterService; private EdgeStatsCounterService statsCounterService;
private EdgeStatsService edgeStatsService; private EdgeStatsService edgeStatsService;
@Captor
private ArgumentCaptor<List<TsKvEntry>> captor;
private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
private final EdgeId edgeId = new EdgeId(UUID.randomUUID()); private final EdgeId edgeId = new EdgeId(UUID.randomUUID());
@BeforeEach @BeforeEach
void setUp() { void setUp() {
edgeStatsService = new EdgeStatsService( edgeStatsService = createEdgeStatsService(Optional.empty());
}
private EdgeStatsService createEdgeStatsService(Optional<KafkaAdmin> kafkaAdmin) {
EdgeStatsService service = new EdgeStatsService(
tsService, tsService,
statsCounterService, statsCounterService,
topicService, topicService,
Optional.empty() kafkaAdmin
); );
ReflectionTestUtils.setField(service, "edgesStatsTtlDays", TTL_DAYS);
ReflectionTestUtils.setField(edgeStatsService, "edgesStatsTtlDays", 30); ReflectionTestUtils.setField(service, "reportIntervalMillis", REPORT_INTERVAL_MILLIS);
ReflectionTestUtils.setField(edgeStatsService, "reportIntervalMillis", 600_000L); return service;
} }
@Test @Test
public void testReportStatsSavesTelemetry() { public void testReportStatsSavesTelemetry() {
// given // GIVEN
setupCounters();
// WHEN
edgeStatsService.reportStats();
// THEN
Map<String, Long> counters = verifyCounters();
Assertions.assertEquals(EXPECTED_MSGS_LAG, counters.get(DOWNLINK_MSGS_LAG.getKey()).longValue());
}
@Test
public void testReportStatsWithKafkaLag() {
// GIVEN
setupCounters();
setupKafkaLag();
// WHEN
edgeStatsService.reportStats();
// THEN
Map<String, Long> valuesByKey = verifyCounters();
Assertions.assertEquals(EXPECTED_MSGS_KAFKA_LAG, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey()));
}
private void setupCounters() {
MsgCounters counters = new MsgCounters(tenantId); MsgCounters counters = new MsgCounters(tenantId);
counters.getMsgsAdded().set(5); counters.getMsgsAdded().set(EXPECTED_MSGS_ADDED);
counters.getMsgsPushed().set(3); counters.getMsgsPushed().set(EXPECTED_MSGS_PUSHED);
counters.getMsgsPermanentlyFailed().set(1); counters.getMsgsPermanentlyFailed().set(EXPECTED_MSGS_PERMANENTLY_FAILED);
counters.getMsgsTmpFailed().set(0); counters.getMsgsTmpFailed().set(EXPECTED_MSGS_TMP_FAILED);
counters.getMsgsLag().set(10); counters.getMsgsLag().set(EXPECTED_MSGS_LAG);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>(); ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters); countersByEdge.put(edgeId, counters);
when(statsCounterService.getCounterByEdge()).thenReturn(countersByEdge); when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge);
ArgumentCaptor<List<TsKvEntry>> captor = ArgumentCaptor.forClass(List.class);
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));
}
// when private Map<String, Long> verifyCounters() {
edgeStatsService.reportStats(); verify(tsService, times(1)).save(eq(tenantId), eq(edgeId), anyList(), anyLong());
verify(statsCounterService, times(1)).clear(edgeId);
// then
List<TsKvEntry> entries = captor.getValue(); List<TsKvEntry> entries = captor.getValue();
Assertions.assertEquals(5, entries.size()); Assertions.assertEquals(5, entries.size());
Map<String, Long> valuesByKey = entries.stream() Map<String, Long> valuesByKey = entries.stream()
.collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L))); .collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L)));
Assertions.assertEquals(5L, valuesByKey.get(DOWNLINK_MSGS_ADDED.getKey()).longValue()); Assertions.assertEquals(EXPECTED_MSGS_ADDED, valuesByKey.get(DOWNLINK_MSGS_ADDED.getKey()).longValue());
Assertions.assertEquals(3L, valuesByKey.get(DOWNLINK_MSGS_PUSHED.getKey()).longValue()); Assertions.assertEquals(EXPECTED_MSGS_PUSHED, valuesByKey.get(DOWNLINK_MSGS_PUSHED.getKey()).longValue());
Assertions.assertEquals(1L, valuesByKey.get(DOWNLINK_MSGS_PERMANENTLY_FAILED.getKey()).longValue()); Assertions.assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, valuesByKey.get(DOWNLINK_MSGS_PERMANENTLY_FAILED.getKey()).longValue());
Assertions.assertEquals(0L, valuesByKey.get(DOWNLINK_MSGS_TMP_FAILED.getKey()).longValue()); Assertions.assertEquals(EXPECTED_MSGS_TMP_FAILED, valuesByKey.get(DOWNLINK_MSGS_TMP_FAILED.getKey()).longValue());
Assertions.assertEquals(10L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey()).longValue()); return valuesByKey;
verify(statsCounterService).clear(edgeId);
} }
@Test private void setupKafkaLag() {
public void testReportStatsWithKafkaLag() {
// given
MsgCounters counters = new MsgCounters(tenantId);
counters.getMsgsAdded().set(2);
counters.getMsgsPushed().set(2);
counters.getMsgsPermanentlyFailed().set(0);
counters.getMsgsTmpFailed().set(1);
counters.getMsgsLag().set(0);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters);
// mocks
when(statsCounterService.getCounterByEdge()).thenReturn(countersByEdge);
String topic = "edge-topic"; String topic = "edge-topic";
TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false); TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false);
when(topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId)).thenReturn(partitionInfo); when(topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId)).thenReturn(partitionInfo);
KafkaAdmin kafkaAdmin = mock(KafkaAdmin.class); KafkaAdmin kafkaAdmin = mock(KafkaAdmin.class);
when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic))) when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic)))
.thenReturn(Map.of(topic, 15L)); .thenReturn(Map.of(topic, EXPECTED_MSGS_KAFKA_LAG));
ArgumentCaptor<List<TsKvEntry>> captor = ArgumentCaptor.forClass(List.class);
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));
edgeStatsService = new EdgeStatsService(
tsService,
statsCounterService,
topicService,
Optional.of(kafkaAdmin)
);
ReflectionTestUtils.setField(edgeStatsService, "edgesStatsTtlDays", 30);
ReflectionTestUtils.setField(edgeStatsService, "reportIntervalMillis", 600_000L);
// when
edgeStatsService.reportStats();
// then
List<TsKvEntry> entries = captor.getValue();
Map<String, Long> valuesByKey = entries.stream()
.collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L)));
Assertions.assertEquals(15L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey())); edgeStatsService = createEdgeStatsService(Optional.of(kafkaAdmin));
verify(statsCounterService).clear(edgeId);
} }
} }

1
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java

@ -41,4 +41,5 @@ public interface EdgeRpcClient {
void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg);
int getServerMaxInboundMessageSize(); int getServerMaxInboundMessageSize();
} }

3
dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java

@ -68,7 +68,8 @@ public class PostgresEdgeEventService extends BaseEdgeEventService {
Futures.addCallback(saveFuture, new FutureCallback<>() { Futures.addCallback(saveFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(Void result) { public void onSuccess(Void result) {
statsCounterService.ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
eventPublisher.publishEvent(SaveEntityEvent.builder() eventPublisher.publishEvent(SaveEntityEvent.builder()
.tenantId(edgeEvent.getTenantId()) .tenantId(edgeEvent.getTenantId())
.entityId(edgeEvent.getEdgeId()) .entityId(edgeEvent.getEdgeId())

15
dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java

@ -24,13 +24,13 @@ import org.thingsboard.server.common.data.id.TenantId;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true", matchIfMissing = false) @ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true")
@Service @Service
@Slf4j @Slf4j
@Getter @Getter
public class EdgeStatsCounterService { public class EdgeStatsCounterService {
private final ConcurrentHashMap<EdgeId, MsgCounters> counterByEdge = new ConcurrentHashMap<>(); private final ConcurrentHashMap<EdgeId, MsgCounters> msgCountersByEdge = new ConcurrentHashMap<>();
public void recordEvent(EdgeStatsKey type, TenantId tenantId, EdgeId edgeId, long value) { public void recordEvent(EdgeStatsKey type, TenantId tenantId, EdgeId edgeId, long value) {
MsgCounters counters = getOrCreateCounters(tenantId, edgeId); MsgCounters counters = getOrCreateCounters(tenantId, edgeId);
@ -39,19 +39,16 @@ public class EdgeStatsCounterService {
case DOWNLINK_MSGS_PUSHED -> counters.getMsgsPushed().addAndGet(value); case DOWNLINK_MSGS_PUSHED -> counters.getMsgsPushed().addAndGet(value);
case DOWNLINK_MSGS_PERMANENTLY_FAILED -> counters.getMsgsPermanentlyFailed().addAndGet(value); case DOWNLINK_MSGS_PERMANENTLY_FAILED -> counters.getMsgsPermanentlyFailed().addAndGet(value);
case DOWNLINK_MSGS_TMP_FAILED -> counters.getMsgsTmpFailed().addAndGet(value); case DOWNLINK_MSGS_TMP_FAILED -> counters.getMsgsTmpFailed().addAndGet(value);
case DOWNLINK_MSGS_LAG -> counters.getMsgsLag().set(value);
} }
} }
public void setDownlinkMsgsLag(TenantId tenantId, EdgeId edgeId, long value) { public MsgCounters getOrCreateCounters(TenantId tenantId, EdgeId edgeId) {
getOrCreateCounters(tenantId, edgeId).getMsgsLag().set(value); return msgCountersByEdge.computeIfAbsent(edgeId, id -> new MsgCounters(tenantId));
} }
public void clear(EdgeId edgeId) { public void clear(EdgeId edgeId) {
counterByEdge.remove(edgeId); msgCountersByEdge.remove(edgeId);
}
public MsgCounters getOrCreateCounters(TenantId tenantId, EdgeId edgeId) {
return counterByEdge.computeIfAbsent(edgeId, id -> new MsgCounters(tenantId));
} }
} }

Loading…
Cancel
Save