From 8ac80476ad2f3768d4f3ecefadd401da607b70a1 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 27 Nov 2023 23:13:03 +0100 Subject: [PATCH] added encode/decode time stats and refactoring --- .../cache/RedisTbTransactionalCache.java | 4 ++++ .../server/common/data/FstStatsService.java | 4 ++++ .../queue/util/ProtoWithFSTService.java | 8 +++++++- .../common/stats/FstStatsServiceImpl.java | 20 +++++++++++++++++-- 4 files changed, 33 insertions(+), 3 deletions(-) diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java b/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java index c2dd66dde0..01a18ad7ad 100644 --- a/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java +++ b/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java @@ -84,8 +84,10 @@ public abstract class RedisTbTransactionalCache clazz); + void recordEncodeTime(Class clazz, long startTime); + + void recordDecodeTime(Class clazz, long startTime); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java index e35c86e070..204502068c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java @@ -36,8 +36,12 @@ public class ProtoWithFSTService implements DataDecodingEncodingService { @Override public Optional decode(byte[] byteArray) { try { + long startTime = System.nanoTime(); Optional optional = Optional.ofNullable(FSTUtils.decode(byteArray)); - optional.ifPresent(obj -> fstStatsService.incrementDecode(obj.getClass())); + optional.ifPresent(obj -> { + fstStatsService.recordDecodeTime(obj.getClass(), startTime); + fstStatsService.incrementDecode(obj.getClass()); + }); return optional; } catch (IllegalArgumentException e) { log.error("Error during deserialization message, [{}]", e.getMessage()); @@ -48,7 +52,9 @@ public class ProtoWithFSTService implements DataDecodingEncodingService { @Override public byte[] encode(T msq) { + long startTime = System.nanoTime(); var bytes = FSTUtils.encode(msq); + fstStatsService.recordEncodeTime(msq.getClass(), startTime); fstStatsService.incrementEncode(msq.getClass()); return bytes; } diff --git a/common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java b/common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java index 5d78b8c36a..b5cb25f2dc 100644 --- a/common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java +++ b/common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java @@ -15,28 +15,44 @@ */ package org.thingsboard.server.common.stats; +import io.micrometer.core.instrument.Timer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.FstStatsService; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; @Service public class FstStatsServiceImpl implements FstStatsService { private final ConcurrentHashMap encodeCounters = new ConcurrentHashMap<>(); private final ConcurrentHashMap decodeCounters = new ConcurrentHashMap<>(); + private final ConcurrentHashMap encodeTimers = new ConcurrentHashMap<>(); + private final ConcurrentHashMap decodeTimer = new ConcurrentHashMap<>(); @Autowired private StatsFactory statsFactory; @Override public void incrementEncode(Class clazz) { - encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstEncode", key)).increment(); + encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fst_encode", key)).increment(); } @Override public void incrementDecode(Class clazz) { - decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstDecode", key)).increment(); + decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fst_decode", key)).increment(); + } + + @Override + public void recordEncodeTime(Class clazz, long startTime) { + encodeTimers.computeIfAbsent(clazz.getSimpleName(), + key -> statsFactory.createTimer("fst_encode_time", "statsName", key)).record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS); + } + + @Override + public void recordDecodeTime(Class clazz, long startTime) { + decodeTimer.computeIfAbsent(clazz.getSimpleName(), + key -> statsFactory.createTimer("fst_decode_time", "statsName", key)).record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS); } }