Browse Source

added encode/decode time stats and refactoring

pull/9708/head
YevhenBondarenko 3 years ago
parent
commit
8ac80476ad
  1. 4
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java
  2. 4
      common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java
  3. 8
      common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java
  4. 20
      common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java

4
common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java

@ -84,8 +84,10 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
} else if (Arrays.equals(rawValue, BINARY_NULL_VALUE)) { } else if (Arrays.equals(rawValue, BINARY_NULL_VALUE)) {
return SimpleTbCacheValueWrapper.empty(); return SimpleTbCacheValueWrapper.empty();
} else { } else {
long startTime = System.nanoTime();
V value = valueSerializer.deserialize(key, rawValue); V value = valueSerializer.deserialize(key, rawValue);
if (value != null) { if (value != null) {
fstStatsService.recordDecodeTime(value.getClass(), startTime);
fstStatsService.incrementDecode(value.getClass()); fstStatsService.incrementDecode(value.getClass());
} }
return SimpleTbCacheValueWrapper.wrap(value); return SimpleTbCacheValueWrapper.wrap(value);
@ -198,7 +200,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return BINARY_NULL_VALUE; return BINARY_NULL_VALUE;
} else { } else {
try { try {
long startTime = System.nanoTime();
var bytes = valueSerializer.serialize(value); var bytes = valueSerializer.serialize(value);
fstStatsService.recordEncodeTime(value.getClass(), startTime);
fstStatsService.incrementEncode(value.getClass()); fstStatsService.incrementEncode(value.getClass());
return bytes; return bytes;
} catch (Exception e) { } catch (Exception e) {

4
common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java

@ -21,4 +21,8 @@ public interface FstStatsService {
void incrementDecode(Class<?> clazz); void incrementDecode(Class<?> clazz);
void recordEncodeTime(Class<?> clazz, long startTime);
void recordDecodeTime(Class<?> clazz, long startTime);
} }

8
common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java

@ -36,8 +36,12 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override @Override
public <T> Optional<T> decode(byte[] byteArray) { public <T> Optional<T> decode(byte[] byteArray) {
try { try {
long startTime = System.nanoTime();
Optional<T> optional = Optional.ofNullable(FSTUtils.decode(byteArray)); Optional<T> 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; return optional;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
log.error("Error during deserialization message, [{}]", e.getMessage()); log.error("Error during deserialization message, [{}]", e.getMessage());
@ -48,7 +52,9 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override @Override
public <T> byte[] encode(T msq) { public <T> byte[] encode(T msq) {
long startTime = System.nanoTime();
var bytes = FSTUtils.encode(msq); var bytes = FSTUtils.encode(msq);
fstStatsService.recordEncodeTime(msq.getClass(), startTime);
fstStatsService.incrementEncode(msq.getClass()); fstStatsService.incrementEncode(msq.getClass());
return bytes; return bytes;
} }

20
common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java

@ -15,28 +15,44 @@
*/ */
package org.thingsboard.server.common.stats; package org.thingsboard.server.common.stats;
import io.micrometer.core.instrument.Timer;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FstStatsService; import org.thingsboard.server.common.data.FstStatsService;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
@Service @Service
public class FstStatsServiceImpl implements FstStatsService { public class FstStatsServiceImpl implements FstStatsService {
private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Timer> encodeTimers = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Timer> decodeTimer = new ConcurrentHashMap<>();
@Autowired @Autowired
private StatsFactory statsFactory; private StatsFactory statsFactory;
@Override @Override
public void incrementEncode(Class<?> clazz) { 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 @Override
public void incrementDecode(Class<?> clazz) { 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);
} }
} }

Loading…
Cancel
Save