Browse Source

fst stats in Grafana

pull/9678/head
YevhenBondarenko 3 years ago
parent
commit
d6b56177a2
  1. 6
      application/src/main/resources/thingsboard.yml
  2. 12
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java
  3. 122
      common/data/src/main/java/org/thingsboard/server/common/data/FSTUtils.java
  4. 15
      common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java
  5. 13
      common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java
  6. 42
      common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java

6
application/src/main/resources/thingsboard.yml

@ -1680,8 +1680,4 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
fst:
stats:
printInterval: "${FST_STATS_PRINT_INTERVAL:600000}" # 10 min by default
include: '${METRICS_ENDPOINTS_EXPOSE:info}'

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

@ -17,6 +17,7 @@ package org.thingsboard.server.cache;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cache.support.NullValue;
import org.springframework.data.redis.connection.RedisConnection;
import org.springframework.data.redis.connection.RedisConnectionFactory;
@ -27,6 +28,7 @@ import org.springframework.data.redis.connection.jedis.JedisConnectionFactory;
import org.springframework.data.redis.core.types.Expiration;
import org.springframework.data.redis.serializer.RedisSerializer;
import org.springframework.data.redis.serializer.StringRedisSerializer;
import org.thingsboard.server.common.data.FstStatsService;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.util.JedisClusterCRC16;
@ -44,6 +46,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
private static final byte[] BINARY_NULL_VALUE = RedisSerializer.java().serialize(NullValue.INSTANCE);
static final JedisPool MOCK_POOL = new JedisPool(); //non-null pool required for JedisConnection to trigger closing jedis connection
@Autowired
private FstStatsService fstStatsService;
@Getter
private final String cacheName;
private final JedisConnectionFactory connectionFactory;
@ -80,6 +85,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return SimpleTbCacheValueWrapper.empty();
} else {
V value = valueSerializer.deserialize(key, rawValue);
if (value != null) {
fstStatsService.incrementDecode(value.getClass());
}
return SimpleTbCacheValueWrapper.wrap(value);
}
}
@ -190,7 +198,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return BINARY_NULL_VALUE;
} else {
try {
return valueSerializer.serialize(value);
var bytes = valueSerializer.serialize(value);
fstStatsService.incrementEncode(value.getClass());
return bytes;
} catch (Exception e) {
log.warn("Failed to serialize the cache value: {}", value, e);
throw new RuntimeException(e);

122
common/data/src/main/java/org/thingsboard/server/common/data/FSTUtils.java

@ -15,137 +15,21 @@
*/
package org.thingsboard.server.common.data;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.nustaq.serialization.FSTConfiguration;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
@Slf4j
public class FSTUtils {
public static final FSTConfiguration CONFIG = FSTConfiguration.createDefaultConfiguration();
private static final ConcurrentHashMap<String, Stats> STATS = new ConcurrentHashMap<>();
@SuppressWarnings("unchecked")
public static <T> T decode(byte[] byteArray) {
long startTime = System.nanoTime();
T result = byteArray != null && byteArray.length > 0 ? (T) CONFIG.asObject(byteArray) : null;
long endTime = System.nanoTime();
if (log.isDebugEnabled() && result != null) {
String className = result.getClass().getSimpleName();
STATS.computeIfAbsent(className, k -> new Stats()).incrementDecode(endTime - startTime);
}
return result;
}
public static <T> byte[] encode(T msg) {
long startTime = System.nanoTime();
byte[] result = CONFIG.asByteArray(msg);
long endTime = System.nanoTime();
if (log.isDebugEnabled() && msg != null) {
String className = msg.getClass().getSimpleName();
STATS.computeIfAbsent(className, k -> new Stats()).incrementEncode(endTime - startTime);
}
return result;
}
public static void printStats() {
if (log.isDebugEnabled()) {
List<String> topDecode = STATS.entrySet().stream()
.filter(e -> e.getValue().getAvgDecodeTime() > 0)
.sorted((e1, e2) -> Long.compare(e2.getValue().getAvgDecodeTime(), e1.getValue().getAvgDecodeTime()))
.limit(5)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
List<String> topEncode = STATS.entrySet().stream()
.filter(e -> e.getValue().getAvgEncodeTime() > 0)
.sorted((e1, e2) -> Long.compare(e2.getValue().getAvgEncodeTime(), e1.getValue().getAvgEncodeTime()))
.limit(5)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
List<String> topDecodeCount = STATS.entrySet().stream()
.filter(e -> e.getValue().getDecodeCount().get() > 0)
.sorted((e1, e2) -> Long.compare(e2.getValue().getDecodeCount().get(), e1.getValue().getDecodeCount().get()))
.limit(5)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
List<String> topEncodeCount = STATS.entrySet().stream()
.filter(e -> e.getValue().getEncodeCount().get() > 0)
.sorted((e1, e2) -> Long.compare(e2.getValue().getEncodeCount().get(), e1.getValue().getEncodeCount().get()))
.limit(5)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
for (Map.Entry<String, Stats> entry : STATS.entrySet()) {
Stats stats = entry.getValue();
if (stats.isNotEmpty()) {
log.debug("[FST stats] [{}] {}", entry.getKey(), stats);
stats.reset();
}
}
log.debug("[FST stats] Top 5 slowest 'decode' {}", topDecode);
log.debug("[FST stats] Top 5 slowest 'encode' {}", topEncode);
log.debug("[FST stats] Top 5 'decode' count {}", topDecodeCount);
log.debug("[FST stats] Top 5 'encode' count {}", topEncodeCount);
}
return byteArray != null && byteArray.length > 0 ? (T) CONFIG.asObject(byteArray) : null;
}
@Data
private static class Stats {
private final AtomicLong encodeCount = new AtomicLong();
private final AtomicLong decodeCount = new AtomicLong();
private final AtomicLong totalEncodeTime = new AtomicLong();
private final AtomicLong totalDecodeTime = new AtomicLong();
private void incrementEncode(long time) {
encodeCount.incrementAndGet();
totalEncodeTime.addAndGet(time);
}
private void incrementDecode(long time) {
decodeCount.incrementAndGet();
totalDecodeTime.addAndGet(time);
}
private boolean isNotEmpty() {
return encodeCount.get() > 0 || decodeCount.get() > 0;
}
private long getAvgEncodeTime() {
long count = encodeCount.get();
return count > 0 ? totalEncodeTime.get() / count : 0;
}
private long getAvgDecodeTime() {
long count = decodeCount.get();
return count > 0 ? totalDecodeTime.get() / count : 0;
}
private void reset() {
encodeCount.set(0);
decodeCount.set(0);
totalEncodeTime.set(0);
totalDecodeTime.set(0);
}
@Override
public String toString() {
return String.format("decodeCount [%d] avgDecodeTime [%d] encodedCount [%d] avgEncodeTime [%d]", decodeCount.get(), getAvgDecodeTime(), encodeCount.get(), getAvgEncodeTime());
}
public static <T> byte[] encode(T msq) {
return CONFIG.asByteArray(msq);
}
}

15
application/src/main/java/org/thingsboard/server/service/stats/FstStats.java → common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java

@ -13,17 +13,12 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.stats;
package org.thingsboard.server.common.data;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FSTUtils;
public interface FstStatsService {
@Service
public class FstStats {
void incrementEncode(Class<?> clazz);
void incrementDecode(Class<?> clazz);
@Scheduled(initialDelayString = "${fst.stats.printInterval}", fixedDelayString = "${fst.stats.printInterval}")
public void printStats() {
FSTUtils.printStats();
}
}

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

@ -17,8 +17,10 @@ package org.thingsboard.server.queue.util;
import lombok.extern.slf4j.Slf4j;
import org.nustaq.serialization.FSTConfiguration;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FSTUtils;
import org.thingsboard.server.common.data.FstStatsService;
import java.util.Optional;
@ -26,12 +28,17 @@ import java.util.Optional;
@Service
public class ProtoWithFSTService implements DataDecodingEncodingService {
@Autowired
private FstStatsService fstStatsService;
public static final FSTConfiguration CONFIG = FSTConfiguration.createDefaultConfiguration();
@Override
public <T> Optional<T> decode(byte[] byteArray) {
try {
return Optional.ofNullable(FSTUtils.decode(byteArray));
Optional<T> optional = Optional.ofNullable(FSTUtils.decode(byteArray));
optional.ifPresent(obj -> fstStatsService.incrementDecode(obj.getClass()));
return optional;
} catch (IllegalArgumentException e) {
log.error("Error during deserialization message, [{}]", e.getMessage());
return Optional.empty();
@ -41,7 +48,9 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override
public <T> byte[] encode(T msq) {
return FSTUtils.encode(msq);
var bytes = FSTUtils.encode(msq);
fstStatsService.incrementEncode(msq.getClass());
return bytes;
}

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

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2023 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.common.stats;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FstStatsService;
import java.util.concurrent.ConcurrentHashMap;
@Service
public class FstStatsServiceImpl implements FstStatsService {
private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>();
@Autowired
private StatsFactory statsFactory;
@Override
public void incrementEncode(Class<?> clazz) {
encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstEncode", key)).increment();
}
@Override
public void incrementDecode(Class<?> clazz) {
decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstDecode", key)).increment();
}
}
Loading…
Cancel
Save