From e8b164a03ad2e07fa3c8ecc61daf14e70d0ed7dc Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 9 Dec 2025 16:56:32 +0200 Subject: [PATCH] Refactoring. Remove EdgeStats - uplinkRate not used --- .../service/edge/EdgeContextComponent.java | 1 - .../service/edge/rpc/EdgeGrpcSession.java | 6 ++-- .../service/edge/stats/EdgeStatsService.java | 14 ++++---- .../server/edge/EdgeStatsIntegrationTest.java | 12 +++---- .../server/service/edge/EdgeStatsTest.java | 25 +++++++------- .../server/dao/edge/stats/EdgeStats.java | 34 ------------------- .../edge/stats/EdgeStatsCounterService.java | 13 ++++--- 7 files changed, 34 insertions(+), 71 deletions(-) delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index ec4a43310b..d523a9176c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -206,7 +206,6 @@ public class EdgeContextComponent { @Autowired private Optional statsCounterService; - // processors @Autowired private AlarmProcessor alarmProcessor; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index c0e76c32e4..83639573bb 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/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 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 { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java index 697f7b6216..66be8fe32a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java +++ b/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 statsByEdgeSnapshot = new HashMap<>(statsCounterService.getStatsByEdge()); + Map countersByEdgeSnapshot = new HashMap<>(statsCounterService.getMsgCountersByEdge()); boolean isKafkaStats = kafkaAdmin.isPresent(); - Map lagByEdgeId = isKafkaStats ? getLagByEdgeId(statsByEdgeSnapshot) : Collections.emptyMap(); - statsByEdgeSnapshot.forEach((edgeId, edgeStats) -> { - MsgCounters counters = edgeStats.getMsgCounters(); + Map 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 getLagByEdgeId(Map edgeStatsByEdge) { - Map edgeToTopicMap = edgeStatsByEdge.entrySet().stream() + private Map getLagByEdgeId(Map countersByEdge) { + Map 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 lagByTopic = kafkaAdmin.get().getTotalLagForGroupsBulk(new HashSet<>(edgeToTopicMap.values())); diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java index b993e4f564..ed98ebe16d 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java +++ b/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 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 toMap(Map.Entry>... suppliers) { + private Map toMap(Map.Entry>... suppliers) { return Arrays.stream(suppliers) .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get())); } diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java index fb24cdb5bb..0b0f28e923 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java +++ b/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> 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 edgeStatsByEdge = new ConcurrentHashMap<>(); - edgeStatsByEdge.put(edgeId, edgeStats); + ConcurrentHashMap countersByEdge = new ConcurrentHashMap<>(); + countersByEdge.put(edgeId, counters); - when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge); + when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge); - ArgumentCaptor> 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 edgeStatsByEdge = new ConcurrentHashMap<>(); - edgeStatsByEdge.put(edgeId, edgeStats); + ConcurrentHashMap 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> captor = ArgumentCaptor.forClass((Class) List.class); when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java deleted file mode 100644 index 36d87691cc..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java +++ /dev/null @@ -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 uplinkRate; - - public EdgeStats(TenantId tenantId) { - this.msgCounters = new MsgCounters(tenantId); - this.uplinkRate = new ConcurrentLinkedQueue<>(); - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java index c537fb75be..ea8b85a287 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java +++ b/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 statsByEdge = new ConcurrentHashMap<>(); + private final ConcurrentHashMap 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); } }