Browse Source

Merge pull request #11607 from YevhenBondarenko/feature/gateway-latency

[MQTT transport] Gateway latency metrics
pull/11668/head
Andrew Shvayka 2 years ago
committed by GitHub
parent
commit
05a813a2c8
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/main/resources/thingsboard.yml
  2. 120
      application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
  3. 19
      common/message/src/main/java/org/thingsboard/server/common/msg/gateway/metrics/GatewayMetadata.java
  4. 45
      common/proto/src/main/java/org/thingsboard/server/common/adaptor/JsonConverter.java
  5. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java
  6. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  7. 2
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java
  8. 26
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TbMqttTransportComponent.java
  9. 113
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/GatewayMetricsService.java
  10. 123
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/metrics/GatewayMetricsState.java
  11. 13
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  12. 18
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
  13. 2
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml

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

@ -1020,6 +1020,8 @@ transport:
# MQTT disconnect timeout in milliseconds. The time to wait for the client to disconnect after the server sends a disconnect message.
disconnect_timeout: "${MQTT_DISCONNECT_TIMEOUT:1000}"
msg_queue_size_per_device_limit: "${MQTT_MSG_QUEUE_SIZE_PER_DEVICE_LIMIT:100}" # messages await in the queue before the device connected state. This limit works on the low level before TenantProfileLimits mechanism
# Interval of periodic report of the gateway metrics
gateway_metrics_report_interval_sec: "${MQTT_GATEWAY_METRICS_REPORT_INTERVAL_SEC:60}"
netty:
# Netty leak detector level
leak_detector_level: "${NETTY_LEAK_DETECTOR_LVL:DISABLED}"

120
application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java

@ -20,22 +20,34 @@ import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.mockito.ArgumentCaptor;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import org.thingsboard.server.transport.mqtt.gateway.GatewayMetricsService;
import org.thingsboard.server.common.msg.gateway.metrics.GatewayMetadata;
import org.thingsboard.server.transport.mqtt.gateway.metrics.GatewayMetricsState;
import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Random;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.mockito.Mockito.any;
import static org.mockito.Mockito.eq;
import static org.mockito.Mockito.verify;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_TELEMETRY_SHORT_JSON_TOPIC;
@ -43,6 +55,7 @@ import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVIC
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_TELEMETRY_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_CONNECT_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_TELEMETRY_TOPIC;
import static org.thingsboard.server.transport.mqtt.gateway.GatewayMetricsService.GATEWAY_METRICS;
@Slf4j
public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqttIntegrationTest {
@ -53,6 +66,9 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
protected static final String MALFORMED_JSON_PAYLOAD = "{\"key1\":, \"key2\":true, \"key3\": 3.0, \"key4\": 4," +
" \"key5\": {\"someNumber\": 42, \"someArray\": [1,2,3], \"someNestedObject\": {\"key\": \"value\"}}}";
@SpyBean
GatewayMetricsService gatewayMetricsService;
@Before
public void beforeTest() throws Exception {
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
@ -106,13 +122,98 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
String deviceName = "Device A";
Device device = doExecuteWithRetriesAndInterval(() -> doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class),
20,
100);
20,
100);
assertNotNull(device);
client.disconnect();
}
@Test
public void testPushMetricsGateway() throws Exception {
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
.gatewayName("Test metrics gateway")
.build();
processBeforeTest(configProperties);
Map<String, List<Long>> gwLatencies = new HashMap<>();
Map<String, List<Long>> transportLatencies = new HashMap<>();
publishLatency(gwLatencies, transportLatencies, 5);
gatewayMetricsService.reportMetrics();
List<String> actualKeys = getActualKeysList(savedGateway.getId(), List.of(GATEWAY_METRICS));
assertEquals(GATEWAY_METRICS, actualKeys.get(0));
String telemetryUrl = String.format("/api/plugins/telemetry/DEVICE/%s/values/timeseries?startTs=%d&endTs=%d&keys=%s", savedGateway.getId(), 0, System.currentTimeMillis(), GATEWAY_METRICS);
Map<String, List<Map<String, Object>>> gatewayTelemetry = doGetAsyncTyped(telemetryUrl, new TypeReference<>() {});
Map<String, Object> latencyCheckTelemetry = gatewayTelemetry.get(GATEWAY_METRICS).get(0);
Map<String, GatewayMetricsState.ConnectorMetricsResult> latencyCheckValue = JacksonUtil.fromString((String) latencyCheckTelemetry.get("value"), new TypeReference<>() {});
assertNotNull(latencyCheckValue);
gwLatencies.forEach((connectorName, gwLatencyList) -> {
long avgGwLatency = (long) gwLatencyList.stream().mapToLong(Long::longValue).average().getAsDouble();
long minGwLatency = gwLatencyList.stream().mapToLong(Long::longValue).min().getAsLong();
long maxGwLatency = gwLatencyList.stream().mapToLong(Long::longValue).max().getAsLong();
List<Long> transportLatencyList = transportLatencies.get(connectorName);
assertNotNull(transportLatencyList);
long avgTransportLatency = (long) transportLatencyList.stream().mapToLong(Long::longValue).average().getAsDouble();
long minTransportLatency = transportLatencyList.stream().mapToLong(Long::longValue).min().getAsLong();
long maxTransportLatency = transportLatencyList.stream().mapToLong(Long::longValue).max().getAsLong();
GatewayMetricsState.ConnectorMetricsResult connectorLatencyResult = latencyCheckValue.get(connectorName);
assertNotNull(connectorLatencyResult);
checkConnectorLatencyResult(connectorLatencyResult, avgGwLatency, minGwLatency, maxGwLatency, avgTransportLatency, minTransportLatency, maxTransportLatency);
});
}
private void publishLatency(Map<String, List<Long>> gwLatencies, Map<String, List<Long>> transportLatencies, int n) throws Exception {
Random random = new Random();
for (int i = 0; i < n; i++) {
long publishedTs = System.currentTimeMillis() - 10;
long gatewayLatencyA = random.nextLong(100, 500);
var firstData = new GatewayMetadata("connectorA", publishedTs - gatewayLatencyA, publishedTs);
gwLatencies.computeIfAbsent("connectorA", key -> new ArrayList<>()).add(gatewayLatencyA);
long gatewayLatencyB = random.nextLong(120, 450);
var secondData = new GatewayMetadata("connectorB", publishedTs - gatewayLatencyB, publishedTs);
gwLatencies.computeIfAbsent("connectorB", key -> new ArrayList<>()).add(gatewayLatencyB);
List<String> expectedKeys = Arrays.asList("key1", "key2", "key3", "key4", "key5");
String deviceName1 = "Device A";
String deviceName2 = "Device B";
String firstMetadata = JacksonUtil.writeValueAsString(firstData);
String secondMetadata = JacksonUtil.writeValueAsString(secondData);
String payload = getGatewayTelemetryJsonPayloadWithMetadata(deviceName1, deviceName2, "10000", "20000", firstMetadata, secondMetadata);
processGatewayTelemetryTest(GATEWAY_TELEMETRY_TOPIC, expectedKeys, payload.getBytes(), deviceName1, deviceName2);
ArgumentCaptor<Long> transportReceiveTsCaptorA = ArgumentCaptor.forClass(Long.class);
ArgumentCaptor<Long> transportReceiveTsCaptorB = ArgumentCaptor.forClass(Long.class);
verify(gatewayMetricsService).process(any(), eq(savedGateway.getId()), eq(List.of(firstData)), transportReceiveTsCaptorA.capture());
verify(gatewayMetricsService).process(any(), eq(savedGateway.getId()), eq(List.of(secondData)), transportReceiveTsCaptorB.capture());
Long transportReceiveTsA = transportReceiveTsCaptorA.getValue();
Long transportReceiveTsB = transportReceiveTsCaptorB.getValue();
Long transportLatencyA = transportReceiveTsA - publishedTs;
Long transportLatencyB = transportReceiveTsB - publishedTs;
transportLatencies.computeIfAbsent("connectorA", key -> new ArrayList<>()).add(transportLatencyA);
transportLatencies.computeIfAbsent("connectorB", key -> new ArrayList<>()).add(transportLatencyB);
}
}
private void checkConnectorLatencyResult(GatewayMetricsState.ConnectorMetricsResult result, long avgGwLatency, long minGwLatency, long maxGwLatency,
long avgTransportLatency, long minTransportLatency, long maxTransportLatency) {
assertNotNull(result);
assertEquals(avgGwLatency, result.avgGwLatency());
assertEquals(minGwLatency, result.minGwLatency());
assertEquals(maxGwLatency, result.maxGwLatency());
assertEquals(avgTransportLatency, result.avgTransportLatency());
assertEquals(minTransportLatency, result.minTransportLatency());
assertEquals(maxTransportLatency, result.maxTransportLatency());
}
protected void processJsonPayloadTelemetryTest(String topic, List<String> expectedKeys, byte[] payload, boolean withTs) throws Exception {
processTelemetryTest(topic, expectedKeys, payload, withTs, false);
}
@ -143,7 +244,8 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
long end = System.currentTimeMillis() + 5000;
Map<String, List<Map<String, Object>>> values = null;
while (start <= end) {
values = doGetAsyncTyped(getTelemetryValuesUrl, new TypeReference<>() {});
values = doGetAsyncTyped(getTelemetryValuesUrl, new TypeReference<>() {
});
boolean valid = values.size() == expectedKeys.size();
if (valid) {
for (String key : expectedKeys) {
@ -182,7 +284,7 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
MqttTestClient client = new MqttTestClient();
client.connectAndWait(gatewayAccessToken);
client.publishAndWait(topic, payload);
client.disconnect();
client.disconnectAndWait();
Device firstDevice = doExecuteWithRetriesAndInterval(() -> doGet("/api/tenant/devices?deviceName=" + firstDeviceName, Device.class),
20,
@ -239,6 +341,16 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
return "{\"" + deviceA + "\": " + payload + ", \"" + deviceB + "\": " + payload + "}";
}
protected String getGatewayTelemetryJsonPayloadWithMetadata(String deviceA, String deviceB, String firstTsValue, String secondTsValue, String firstMetadata, String secondMetadata) {
String payloadA = "[{\"ts\": " + firstTsValue + ", \"values\": " + PAYLOAD_VALUES_STR + ", \"metadata\":" + firstMetadata + "}, " +
"{\"ts\": " + secondTsValue + ", \"values\": " + PAYLOAD_VALUES_STR + "}]";
String payloadB = "[{\"ts\": " + firstTsValue + ", \"values\": " + PAYLOAD_VALUES_STR + ", \"metadata\":" + secondMetadata + "}, " +
"{\"ts\": " + secondTsValue + ", \"values\": " + PAYLOAD_VALUES_STR + "}]";
return "{\"" + deviceA + "\": " + payloadA + ", \"" + deviceB + "\": " + payloadB + "}";
}
private String getTelemetryValuesUrl(DeviceId deviceId, Set<String> actualKeySet) {
return "/api/plugins/telemetry/DEVICE/" + deviceId + "/values/timeseries?startTs=0&endTs=25000&keys=" + String.join(",", actualKeySet);
}

19
common/message/src/main/java/org/thingsboard/server/common/msg/gateway/metrics/GatewayMetadata.java

@ -0,0 +1,19 @@
/**
* Copyright © 2016-2024 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.msg.gateway.metrics;
public record GatewayMetadata(String connector, long receivedTs, long publishedTs) {
}

45
common/proto/src/main/java/org/thingsboard/server/common/adaptor/JsonConverter.java

@ -34,6 +34,8 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.gateway.metrics.GatewayMetadata;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg;
@ -83,6 +85,49 @@ public class JsonConverter {
return convertToTelemetryProto(jsonElement, System.currentTimeMillis());
}
public static TbPair<TransportProtos.PostTelemetryMsg, List<GatewayMetadata>> convertToGatewayTelemetry(JsonElement jsonElement, long systemTs) {
List<GatewayMetadata> metadataResult = null;
PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder();
if (jsonElement.isJsonArray()) {
var ja = jsonElement.getAsJsonArray();
for (int i = 0; i < ja.size(); i++) {
var je = ja.get(i);
if (je.isJsonObject()) {
JsonObject jo = je.getAsJsonObject();
JsonElement metadataElem = jo.remove("metadata");
if (metadataElem != null) {
if (metadataResult == null) {
metadataResult = new ArrayList<>();
}
if (metadataElem.isJsonObject()) {
JsonObject metadataObj = metadataElem.getAsJsonObject();
var connector = getAndValidateMetadataElement(metadataObj, "connector").getAsString();
var receivedTs = getAndValidateMetadataElement(metadataObj, "receivedTs").getAsLong();
var publishedTs = getAndValidateMetadataElement(metadataObj, "publishedTs").getAsLong();
metadataResult.add(new GatewayMetadata(connector, receivedTs, publishedTs));
} else {
throw new JsonSyntaxException("Can't parse gateway metadata: " + metadataElem);
}
}
parseObject(systemTs, null, builder, jo);
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + je);
}
}
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + jsonElement);
}
return TbPair.of(builder.build(), metadataResult);
}
private static JsonElement getAndValidateMetadataElement(JsonObject metadata, String elementName) {
var element = metadata.get(elementName);
if (element == null || element.isJsonNull()) {
throw new JsonSyntaxException(String.format("Can't parse gateway element in metadata: [%s][%s]", metadata, elementName));
}
return element;
}
private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder) {
if (jsonElement.isJsonObject()) {
parseObject(systemTs, result, builder, jsonElement.getAsJsonObject());

8
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java

@ -22,12 +22,12 @@ import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor;
import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor;
import org.thingsboard.server.transport.mqtt.gateway.GatewayMetricsService;
import java.net.InetSocketAddress;
import java.util.concurrent.atomic.AtomicInteger;
@ -37,7 +37,7 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
@Slf4j
@Component
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')")
@TbMqttTransportComponent
public class MqttTransportContext extends TransportContext {
@Getter
@ -56,6 +56,10 @@ public class MqttTransportContext extends TransportContext {
@Autowired
private TransportTenantProfileCache tenantProfileCache;
@Getter
@Autowired
private GatewayMetricsService gatewayMetricsService;
@Getter
@Value("${transport.mqtt.netty.max_payload_size}")
private Integer maxPayloadSize;

5
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -1464,7 +1464,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
public void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional<DeviceProfile> deviceProfileOpt) {
deviceSessionCtx.onDeviceUpdate(sessionInfo, device, deviceProfileOpt);
if (gatewaySessionHandler != null) {
gatewaySessionHandler.onDeviceUpdate(sessionInfo, device, deviceProfileOpt);
gatewaySessionHandler.onGatewayUpdate(sessionInfo, device, deviceProfileOpt);
}
}
@ -1473,6 +1473,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
context.onAuthFailure(address);
ChannelHandlerContext ctx = deviceSessionCtx.getChannel();
closeCtx(ctx, MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION);
if (gatewaySessionHandler != null) {
gatewaySessionHandler.onGatewayDelete(deviceId);
}
}
public void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ThingsboardErrorCode result, String errorMsg) {

2
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

@ -39,7 +39,7 @@ import java.net.InetSocketAddress;
* @author Andrew Shvayka
*/
@Service("MqttTransportService")
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')")
@TbMqttTransportComponent
@Slf4j
public class MqttTransportService implements TbTransportService {

26
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TbMqttTransportComponent.java

@ -0,0 +1,26 @@
/**
* Copyright © 2016-2024 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.transport.mqtt;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
@Retention(RetentionPolicy.RUNTIME)
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')")
public @interface TbMqttTransportComponent {
}

113
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/GatewayMetricsService.java

@ -0,0 +1,113 @@
/**
* Copyright © 2016-2024 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.transport.mqtt.gateway;
import jakarta.annotation.PostConstruct;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
import org.thingsboard.server.transport.mqtt.TbMqttTransportComponent;
import org.thingsboard.server.common.msg.gateway.metrics.GatewayMetadata;
import org.thingsboard.server.transport.mqtt.gateway.metrics.GatewayMetricsState;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
@Slf4j
@Service
@TbMqttTransportComponent
public class GatewayMetricsService {
public static final String GATEWAY_METRICS = "gatewayMetrics";
@Value("${transport.mqtt.gateway_metrics_report_interval_sec:60}")
private int metricsReportIntervalSec;
@Autowired
private SchedulerComponent scheduler;
@Autowired
private TransportService transportService;
private Map<DeviceId, GatewayMetricsState> states = new ConcurrentHashMap<>();
@PostConstruct
private void init() {
scheduler.scheduleAtFixedRate(this::reportMetrics, metricsReportIntervalSec, metricsReportIntervalSec, TimeUnit.SECONDS);
}
public void process(TransportProtos.SessionInfoProto sessionInfo, DeviceId gatewayId, List<GatewayMetadata> data, long serverReceiveTs) {
states.computeIfAbsent(gatewayId, k -> new GatewayMetricsState(sessionInfo)).update(data, serverReceiveTs);
}
public void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceId gatewayId) {
var state = states.get(gatewayId);
if (state != null) {
state.updateSessionInfo(sessionInfo);
}
}
public void onDeviceDelete(DeviceId deviceId) {
states.remove(deviceId);
}
public void reportMetrics() {
if (states.isEmpty()) {
return;
}
Map<DeviceId, GatewayMetricsState> statesToReport = states;
states = new ConcurrentHashMap<>();
long ts = System.currentTimeMillis();
statesToReport.forEach((gatewayId, state) -> {
reportMetrics(state, ts);
});
}
private void reportMetrics(GatewayMetricsState state, long ts) {
if (state.isEmpty()) {
return;
}
var result = state.getStateResult();
var kvProto = TransportProtos.KeyValueProto.newBuilder()
.setKey(GATEWAY_METRICS)
.setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV(JacksonUtil.toString(result))
.build();
TransportProtos.TsKvListProto tsKvList = TransportProtos.TsKvListProto.newBuilder()
.setTs(ts)
.addKv(kvProto)
.build();
TransportProtos.PostTelemetryMsg telemetryMsg = TransportProtos.PostTelemetryMsg.newBuilder()
.addTsKvList(tsKvList)
.build();
transportService.process(state.getSessionInfo(), telemetryMsg, TransportServiceCallback.EMPTY);
}
}

123
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/metrics/GatewayMetricsState.java

@ -0,0 +1,123 @@
/**
* Copyright © 2016-2024 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.transport.mqtt.gateway.metrics;
import lombok.Getter;
import org.thingsboard.server.common.msg.gateway.metrics.GatewayMetadata;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
public class GatewayMetricsState {
private final Map<String, ConnectorMetricsState> connectors;
private final Lock updateLock;
@Getter
private volatile TransportProtos.SessionInfoProto sessionInfo;
public GatewayMetricsState(TransportProtos.SessionInfoProto sessionInfo) {
this.connectors = new HashMap<>();
this.updateLock = new ReentrantLock();
this.sessionInfo = sessionInfo;
}
public void updateSessionInfo(TransportProtos.SessionInfoProto sessionInfo) {
this.sessionInfo = sessionInfo;
}
public void update(List<GatewayMetadata> metricsData, long serverReceiveTs) {
updateLock.lock();
try {
metricsData.forEach(data -> {
connectors.computeIfAbsent(data.connector(), k -> new ConnectorMetricsState()).update(data, serverReceiveTs);
});
} finally {
updateLock.unlock();
}
}
public Map<String, ConnectorMetricsResult> getStateResult() {
Map<String, ConnectorMetricsResult> result = new HashMap<>();
updateLock.lock();
try {
connectors.forEach((name, state) -> result.put(name, state.getResult()));
connectors.clear();
} finally {
updateLock.unlock();
}
return result;
}
public boolean isEmpty() {
return connectors.isEmpty();
}
private static class ConnectorMetricsState {
private final AtomicInteger count;
private final AtomicLong gwLatencySum;
private final AtomicLong transportLatencySum;
private volatile long minGwLatency;
private volatile long maxGwLatency;
private volatile long minTransportLatency;
private volatile long maxTransportLatency;
private ConnectorMetricsState() {
this.count = new AtomicInteger(0);
this.gwLatencySum = new AtomicLong(0);
this.transportLatencySum = new AtomicLong(0);
}
private void update(GatewayMetadata metricsData, long serverReceiveTs) {
long gwLatency = metricsData.publishedTs() - metricsData.receivedTs();
long transportLatency = serverReceiveTs - metricsData.publishedTs();
count.incrementAndGet();
gwLatencySum.addAndGet(gwLatency);
transportLatencySum.addAndGet(transportLatency);
if (minGwLatency == 0 || minGwLatency > gwLatency) {
minGwLatency = gwLatency;
}
if (maxGwLatency < gwLatency) {
maxGwLatency = gwLatency;
}
if (minTransportLatency == 0 || minTransportLatency > transportLatency) {
minTransportLatency = transportLatency;
}
if (maxTransportLatency < transportLatency) {
maxTransportLatency = transportLatency;
}
}
private ConnectorMetricsResult getResult() {
long count = this.count.get();
long avgGwLatency = gwLatencySum.get() / count;
long avgTransportLatency = transportLatencySum.get() / count;
return new ConnectorMetricsResult(avgGwLatency, minGwLatency, maxGwLatency, avgTransportLatency, minTransportLatency, maxTransportLatency);
}
}
public record ConnectorMetricsResult(long avgGwLatency, long minGwLatency, long maxGwLatency,
long avgTransportLatency, long minTransportLatency, long maxTransportLatency) {
}
}

13
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java

@ -48,6 +48,8 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.gateway.metrics.GatewayMetadata;
import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
@ -62,6 +64,7 @@ import org.thingsboard.server.transport.mqtt.MqttTransportHandler;
import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor;
import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor;
import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor;
import org.thingsboard.server.transport.mqtt.gateway.GatewayMetricsService;
import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState;
import java.util.ArrayList;
@ -114,6 +117,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
protected final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap;
protected final ChannelHandlerContext channel;
protected final DeviceSessionCtx deviceSessionCtx;
protected final GatewayMetricsService gatewayMetricsService;
@Getter
@Setter
@ -131,6 +135,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
this.mqttQoSMap = deviceSessionCtx.getMqttQoSMap();
this.channel = deviceSessionCtx.getChannel();
this.overwriteDevicesActivity = overwriteDevicesActivity;
this.gatewayMetricsService = deviceSessionCtx.getContext().getGatewayMetricsService();
}
ConcurrentReferenceHashMap<String, Lock> createWeakMap() {
@ -380,7 +385,13 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
private void processPostTelemetryMsg(T deviceCtx, JsonElement msg, String deviceName, int msgId) {
try {
TransportProtos.PostTelemetryMsg postTelemetryMsg = JsonConverter.convertToTelemetryProto(msg.getAsJsonArray());
long systemTs = System.currentTimeMillis();
TbPair<TransportProtos.PostTelemetryMsg, List<GatewayMetadata>> gatewayPayloadPair = JsonConverter.convertToGatewayTelemetry(msg.getAsJsonArray(), systemTs);
TransportProtos.PostTelemetryMsg postTelemetryMsg = gatewayPayloadPair.getFirst();
List<GatewayMetadata> metadata = gatewayPayloadPair.getSecond();
if (!CollectionUtils.isEmpty(metadata)) {
gatewayMetricsService.process(deviceSessionCtx.getSessionInfo(), gateway.getDeviceId(), metadata, systemTs);
}
transportService.process(deviceCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(channel, deviceName, msgId, postTelemetryMsg));
} catch (Throwable e) {
log.warn("[{}][{}][{}] Failed to convert telemetry: [{}]", gateway.getTenantId(), gateway.getDeviceId(), deviceName, msg, e);

18
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java

@ -17,14 +17,21 @@ package org.thingsboard.server.transport.mqtt.session;
import io.netty.buffer.ByteBuf;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.adaptor.AdaptorException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Optional;
import java.util.UUID;
/**
* Created by nickAS21 on 26.12.22
*/
@Slf4j
public class GatewaySessionHandler extends AbstractGatewaySessionHandler<GatewayDeviceSessionContext> {
public GatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId, boolean overwriteDevicesActivity) {
@ -51,7 +58,16 @@ public class GatewaySessionHandler extends AbstractGatewaySessionHandler<Gateway
@Override
protected GatewayDeviceSessionContext newDeviceSessionCtx(GetOrCreateDeviceFromGatewayResponse msg) {
return new GatewayDeviceSessionContext(this, msg.getDeviceInfo(), msg.getDeviceProfile(), mqttQoSMap, transportService);
return new GatewayDeviceSessionContext(this, msg.getDeviceInfo(), msg.getDeviceProfile(), mqttQoSMap, transportService);
}
public void onGatewayUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional<DeviceProfile> deviceProfileOpt) {
this.onDeviceUpdate(sessionInfo, device, deviceProfileOpt);
gatewayMetricsService.onDeviceUpdate(sessionInfo, gateway.getDeviceId());
}
public void onGatewayDelete(DeviceId deviceId) {
gatewayMetricsService.onDeviceDelete(deviceId);
}
}

2
transport/mqtt/src/main/resources/tb-mqtt-transport.yml

@ -146,6 +146,8 @@ transport:
# MQTT disconnect timeout in milliseconds. The time to wait for the client to disconnect after the server sends a disconnect message.
disconnect_timeout: "${MQTT_DISCONNECT_TIMEOUT:1000}"
msg_queue_size_per_device_limit: "${MQTT_MSG_QUEUE_SIZE_PER_DEVICE_LIMIT:100}" # messages await in the queue before device connected state. This limit works on low level before TenantProfileLimits mechanism
# Interval of periodic report of the gateway metrics
gateway_metrics_report_interval_sec: "${MQTT_GATEWAY_METRICS_REPORT_INTERVAL_SEC:60}"
netty:
# Netty leak detector level
leak_detector_level: "${NETTY_LEAK_DETECTOR_LVL:DISABLED}"

Loading…
Cancel
Save