Browse Source

Refactoring. Remove EdgeStats - uplinkRate not used

pull/14549/head
Volodymyr Babak 10 months ago
parent
commit
e8b164a03a
  1. 1
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 14
      application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java
  4. 12
      application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java
  5. 25
      application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java
  6. 34
      dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java
  7. 13
      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
private Optional<EdgeStatsCounterService> statsCounterService;
// processors
@Autowired
private AlarmProcessor alarmProcessor;

6
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 (fetcher instanceof GeneralEdgeEventFetcher) {
long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage());
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, 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());
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())
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg)
.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);
}
} else {

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

@ -31,7 +31,6 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.edge.stats.EdgeStats;
import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService;
import org.thingsboard.server.dao.edge.stats.MsgCounters;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
@ -81,11 +80,10 @@ public class EdgeStatsService {
long now = System.currentTimeMillis();
long ts = now - (now % reportIntervalMillis);
Map<EdgeId, EdgeStats> statsByEdgeSnapshot = new HashMap<>(statsCounterService.getStatsByEdge());
Map<EdgeId, MsgCounters> countersByEdgeSnapshot = new HashMap<>(statsCounterService.getMsgCountersByEdge());
boolean isKafkaStats = kafkaAdmin.isPresent();
Map<EdgeId, Long> lagByEdgeId = isKafkaStats ? getLagByEdgeId(statsByEdgeSnapshot) : Collections.emptyMap();
statsByEdgeSnapshot.forEach((edgeId, edgeStats) -> {
MsgCounters counters = edgeStats.getMsgCounters();
Map<EdgeId, Long> lagByEdgeId = isKafkaStats ? getLagByEdgeId(countersByEdgeSnapshot) : Collections.emptyMap();
countersByEdgeSnapshot.forEach((edgeId, counters) -> {
TenantId tenantId = counters.getTenantId();
if (isKafkaStats) {
@ -104,11 +102,11 @@ public class EdgeStatsService {
});
}
private Map<EdgeId, Long> getLagByEdgeId(Map<EdgeId, EdgeStats> edgeStatsByEdge) {
Map<EdgeId, String> edgeToTopicMap = edgeStatsByEdge.entrySet().stream()
private Map<EdgeId, Long> getLagByEdgeId(Map<EdgeId, MsgCounters> countersByEdge) {
Map<EdgeId, String> edgeToTopicMap = countersByEdge.entrySet().stream()
.collect(Collectors.toMap(
Map.Entry::getKey,
e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getMsgCounters().getTenantId(), e.getKey()).getTopic()
e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getTenantId(), e.getKey()).getTopic()
));
Map<String, Long> lagByTopic = kafkaAdmin.get().getTotalLagForGroupsBulk(new HashSet<>(edgeToTopicMap.values()));

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

@ -179,19 +179,19 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
private void prepareTestData() throws InterruptedException, ExecutionException {
statsCounterService.clear(edge.getId());
// 2 stats message ADDED Device Profile, ASSIGN Device
// 2 stats messages: ADDED Device Profile, ASSIGN Device
edgeImitator.expectMessageAmount(4);
Device device = saveDevice("StatisticDevice", STATISTICS_DEVICE_PROFILE);
doPost("/api/edge/" + edge.getUuidId() + "/device/" + device.getUuidId(), Device.class);
edgeImitator.waitForMessages();
// 1 stats message ASSIGN Asset
// 1 stats message: ASSIGN Asset
edgeImitator.expectMessageAmount(2);
Asset savedAsset = saveAsset("Edge Asset 2");
doPost("/api/edge/" + edge.getUuidId()
+ "/asset/" + savedAsset.getUuidId(), Asset.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// 2 stats message ADDED Customer, ASSIGN Customer
// 2 stats messages: ADDED Customer, ASSIGN Customer
edgeImitator.expectMessageAmount(1);
Customer customer = new Customer();
customer.setTitle("Edge Customer");
@ -204,7 +204,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
+ "/edge/" + edge.getUuidId(), Edge.class);
Assert.assertTrue(edgeImitator.waitForMessages());
//1 stats message Timeseries
// 1 stats message: Timeseries
edgeImitator.expectMessageAmount(1);
String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = JacksonUtil.toJsonNode(timeseriesData);
@ -222,7 +222,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
private void awaitEdgeCountersUpdated() {
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
MsgCounters counters = statsCounterService.getStatsByEdge().get(edge.getId()).getMsgCounters();
MsgCounters counters = statsCounterService.getMsgCountersByEdge().get(edge.getId());
Map<String, Long> actualCounters = toMap(
Map.entry(DOWNLINK_MSGS_ADDED.getKey(), () -> counters.getMsgsAdded().get()),
Map.entry(DOWNLINK_MSGS_PUSHED.getKey(), () -> counters.getMsgsPushed().get()),
@ -234,7 +234,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
}
@SafeVarargs
private final Map<String, Long> toMap(Map.Entry<String, Supplier<Long>>... suppliers) {
private Map<String, Long> toMap(Map.Entry<String, Supplier<Long>>... suppliers) {
return Arrays.stream(suppliers)
.collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get()));
}

25
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.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
@ -29,7 +30,6 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.edge.stats.EdgeStats;
import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService;
import org.thingsboard.server.dao.edge.stats.MsgCounters;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
@ -72,6 +72,9 @@ public class EdgeStatsTest {
private EdgeStatsCounterService statsCounterService;
private EdgeStatsService edgeStatsService;
@Captor
private ArgumentCaptor<List<TsKvEntry>> captor;
private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
private final EdgeId edgeId = new EdgeId(UUID.randomUUID());
@ -94,20 +97,18 @@ public class EdgeStatsTest {
@Test
public void testReportStatsSavesTelemetry() {
EdgeStats edgeStats = new EdgeStats(tenantId);
MsgCounters counters = edgeStats.getMsgCounters();
MsgCounters counters = new MsgCounters(tenantId);
counters.getMsgsAdded().set(5);
counters.getMsgsPushed().set(3);
counters.getMsgsPermanentlyFailed().set(1);
counters.getMsgsTmpFailed().set(0);
counters.getMsgsLag().set(10);
ConcurrentHashMap<EdgeId, EdgeStats> edgeStatsByEdge = new ConcurrentHashMap<>();
edgeStatsByEdge.put(edgeId, edgeStats);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters);
when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge);
when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge);
ArgumentCaptor<List<TsKvEntry>> captor = ArgumentCaptor.forClass((Class) List.class);
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));
@ -131,18 +132,17 @@ public class EdgeStatsTest {
@Test
public void testReportStatsWithKafkaLag() {
EdgeStats edgeStats = new EdgeStats(tenantId);
MsgCounters counters = edgeStats.getMsgCounters();
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, EdgeStats> edgeStatsByEdge = new ConcurrentHashMap<>();
edgeStatsByEdge.put(edgeId, edgeStats);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters);
when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge);
when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge);
String topic = "edge-topic";
TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false);
@ -152,7 +152,6 @@ public class EdgeStatsTest {
when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic)))
.thenReturn(Map.of(topic, 15L));
ArgumentCaptor<List<TsKvEntry>> captor = ArgumentCaptor.forClass((Class) List.class);
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));

34
dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java

@ -1,34 +0,0 @@
/**
* 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.dao.edge.stats;
import lombok.Data;
import org.thingsboard.server.common.data.id.TenantId;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
@Data
public class EdgeStats {
private final MsgCounters msgCounters;
private final Queue<Long> uplinkRate;
public EdgeStats(TenantId tenantId) {
this.msgCounters = new MsgCounters(tenantId);
this.uplinkRate = new ConcurrentLinkedQueue<>();
}
}

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

@ -24,17 +24,16 @@ import org.thingsboard.server.common.data.id.TenantId;
import java.util.concurrent.ConcurrentHashMap;
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true", matchIfMissing = false)
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true")
@Service
@Slf4j
@Getter
public class EdgeStatsCounterService {
private final ConcurrentHashMap<EdgeId, EdgeStats> statsByEdge = new ConcurrentHashMap<>();
private final ConcurrentHashMap<EdgeId, MsgCounters> msgCountersByEdge = new ConcurrentHashMap<>();
public void recordEvent(EdgeStatsKey type, TenantId tenantId, EdgeId edgeId, long value) {
EdgeStats edgeStats = getOrCreateEdgeStats(tenantId, edgeId);
MsgCounters counters = edgeStats.getMsgCounters();
MsgCounters counters = getOrCreateCounters(tenantId, edgeId);
switch (type) {
case DOWNLINK_MSGS_ADDED -> counters.getMsgsAdded().addAndGet(value);
case DOWNLINK_MSGS_PUSHED -> counters.getMsgsPushed().addAndGet(value);
@ -44,12 +43,12 @@ public class EdgeStatsCounterService {
}
}
public EdgeStats getOrCreateEdgeStats(TenantId tenantId, EdgeId edgeId) {
return statsByEdge.computeIfAbsent(edgeId, id -> new EdgeStats(tenantId));
public MsgCounters getOrCreateCounters(TenantId tenantId, EdgeId edgeId) {
return msgCountersByEdge.computeIfAbsent(edgeId, id -> new MsgCounters(tenantId));
}
public void clear(EdgeId edgeId) {
statsByEdge.remove(edgeId);
msgCountersByEdge.remove(edgeId);
}
}

Loading…
Cancel
Save