From cc2446442d7423a7038cb3a95e239e0a4e7c276d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 20 Oct 2020 12:00:35 +0300 Subject: [PATCH] Api Usage Stats flow --- .../DefaultTbApiUsageStateService.java | 105 ++++++++++++++++++ .../apiusage/DefaultTbUsageStatsService.java | 30 ----- ...rvice.java => TbApiUsageStateService.java} | 2 +- .../service/apiusage/TenantApiUsageState.java | 47 ++++++++ .../queue/DefaultTbCoreConsumerService.java | 6 +- .../usagestats/DefaultTbUsageStatsClient.java | 47 ++++---- .../service/DefaultTransportService.java | 36 +++++- 7 files changed, 217 insertions(+), 56 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java delete mode 100644 application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java rename application/src/main/java/org/thingsboard/server/service/apiusage/{TbUsageStatsService.java => TbApiUsageStateService.java} (95%) create mode 100644 application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java new file mode 100644 index 0000000000..8a50c9d62b --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -0,0 +1,105 @@ +/** + * Copyright © 2016-2020 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.service.apiusage; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.ApiUsageRecordKey; +import org.thingsboard.server.common.data.ApiUsageState; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.tools.SchedulerUtils; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; +import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; +import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.util.TbCoreComponent; + +import java.time.LocalDate; +import java.time.ZoneId; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; + +@Slf4j +@TbCoreComponent +@Service +public class DefaultTbApiUsageStateService implements TbApiUsageStateService { + + private final ApiUsageStateService apiUsageStateService; + private final TimeseriesService tsService; + private final ZoneId zoneId; + private final Map tenantStates = new ConcurrentHashMap<>(); + + public DefaultTbApiUsageStateService(ApiUsageStateService apiUsageStateService, TimeseriesService tsService) { + this.apiUsageStateService = apiUsageStateService; + this.tsService = tsService; + this.zoneId = SchedulerUtils.getZoneId("UTC"); + } + + @Override + public void process(TbProtoQueueMsg msg, TbCallback callback) { + ToUsageStatsServiceMsg statsMsg = msg.getValue(); + TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB())); + TenantApiUsageState tenantState = getOrFetchState(tenantId); + long ts = tenantState.getCurrentMonthTs(); + List updatedEntries = new ArrayList<>(ApiUsageRecordKey.values().length); + for (UsageStatsKVProto kvProto : statsMsg.getValuesList()) { + ApiUsageRecordKey recordKey = ApiUsageRecordKey.valueOf(kvProto.getKey()); + long newValue = tenantState.add(recordKey, kvProto.getValue()); + updatedEntries.add(new BasicTsKvEntry(ts, new LongDataEntry(recordKey.name(), newValue))); + } + tsService.save(tenantId, tenantState.getEntityId(), updatedEntries, 0L); + } + + private TenantApiUsageState getOrFetchState(TenantId tenantId) { + TenantApiUsageState tenantState = tenantStates.get(tenantId); + if (tenantState == null) { + long currentMonthTs = LocalDate.now().withDayOfMonth(1).atStartOfDay(zoneId).toInstant().toEpochMilli(); + ApiUsageState dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId); + tenantState = new TenantApiUsageState(currentMonthTs, dbStateEntity.getEntityId()); + try { + List dbValues = tsService.findAllLatest(tenantId, dbStateEntity.getEntityId()).get(); + for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { + TsKvEntry keyEntry = null; + for (TsKvEntry tsKvEntry : dbValues) { + if (tsKvEntry.getKey().equals(key.name())) { + keyEntry = tsKvEntry; + break; + } + } + if (keyEntry != null) { + tenantState.put(key, keyEntry.getLongValue().get()); + } else { + tenantState.put(key, 0L); + } + } + tenantStates.put(tenantId, tenantState); + } catch (InterruptedException | ExecutionException e) { + log.warn("[{}] Failed to fetch api usage state from db.", tenantId, e); + } + } + return tenantState; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java deleted file mode 100644 index e2679d6433..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java +++ /dev/null @@ -1,30 +0,0 @@ -/** - * Copyright © 2016-2020 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.service.apiusage; - -import org.springframework.stereotype.Service; -import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.common.TbProtoQueueMsg; - -@Service -public class DefaultTbUsageStatsService implements TbUsageStatsService { - @Override - public void process(TbProtoQueueMsg msg, TbCallback callback) { - - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TbUsageStatsService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java similarity index 95% rename from application/src/main/java/org/thingsboard/server/service/apiusage/TbUsageStatsService.java rename to application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java index 3c54479eaf..b5aeedeb06 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/TbUsageStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java @@ -19,7 +19,7 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -public interface TbUsageStatsService { +public interface TbApiUsageStateService { void process(TbProtoQueueMsg msg, TbCallback callback); } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java new file mode 100644 index 0000000000..a8be19914b --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java @@ -0,0 +1,47 @@ +/** + * Copyright © 2016-2020 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.service.apiusage; + +import lombok.Getter; +import org.thingsboard.server.common.data.ApiUsageRecordKey; +import org.thingsboard.server.common.data.id.EntityId; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +public class TenantApiUsageState { + + private final Map values = new ConcurrentHashMap<>(); + @Getter + private final EntityId entityId; + @Getter + private volatile long currentMonthTs; + + public TenantApiUsageState(long currentMonthTs, EntityId entityId) { + this.entityId = entityId; + this.currentMonthTs = currentMonthTs; + } + + public void put(ApiUsageRecordKey key, Long value) { + values.put(key, value); + } + + public long add(ApiUsageRecordKey key, long value) { + long result = values.getOrDefault(key, 0L) + value; + values.put(key, result); + return result; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index a389eefd1a..51fce853db 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -62,7 +62,7 @@ import org.thingsboard.server.service.subscription.SubscriptionManagerService; import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; -import org.thingsboard.server.service.apiusage.TbUsageStatsService; +import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -92,7 +92,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> mainConsumer; private final DeviceStateService stateService; - private final TbUsageStatsService statsService; + private final TbApiUsageStateService statsService; private final TbLocalSubscriptionService localSubscriptionService; private final SubscriptionManagerService subscriptionManagerService; private final TbCoreDeviceRpcService tbCoreDeviceRpcService; @@ -106,7 +106,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService[] values = new ConcurrentMap[ApiUsageRecordKey.values().length]; @@ -69,28 +71,33 @@ public class DefaultTbUsageStatsClient implements TbUsageStatsClient { } private void reportStats() { - ConcurrentMap report = new ConcurrentHashMap<>(); + try { + ConcurrentMap report = new ConcurrentHashMap<>(); - for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { - values[key.ordinal()].forEach(((tenantId, atomicLong) -> { - long value = atomicLong.getAndSet(0); - if (value > 0) { - ToUsageStatsServiceMsg.Builder msgBuilder = report.computeIfAbsent(tenantId, id -> { - ToUsageStatsServiceMsg.Builder msg = ToUsageStatsServiceMsg.newBuilder(); - msg.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); - msg.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); - return msg; - }); - msgBuilder.addValues(UsageStatsKVProto.newBuilder().setKey(key.name()).setValue(value).build()); - } + for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { + values[key.ordinal()].forEach(((tenantId, atomicLong) -> { + long value = atomicLong.getAndSet(0); + if (value > 0) { + ToUsageStatsServiceMsg.Builder msgBuilder = report.computeIfAbsent(tenantId, id -> { + ToUsageStatsServiceMsg.Builder msg = ToUsageStatsServiceMsg.newBuilder(); + msg.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + msg.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + return msg; + }); + msgBuilder.addValues(UsageStatsKVProto.newBuilder().setKey(key.name()).setValue(value).build()); + } + })); + } + + report.forEach(((tenantId, builder) -> { + //TODO: figure out how to minimize messages into the queue. Maybe group by 100s of messages? + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId); + msgProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), builder.build()), null); })); + log.info("Report statistics for: {} tenants", report.size()); + } catch (Exception e) { + log.warn("Failed to report statistics: ", e); } - - report.forEach(((tenantId, builder) -> { - //TODO: figure out how to minimize messages into the queue. Maybe group by 100s of messages? - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId); - msgProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), builder.build()), null); - })); } @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 774a437e39..010a3cd899 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -25,6 +25,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.EntityType; @@ -78,6 +79,7 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbTransportQueueFactory; import org.thingsboard.server.queue.scheduler.SchedulerComponent; +import org.thingsboard.server.queue.usagestats.TbUsageStatsClient; import org.thingsboard.server.queue.util.TbTransportComponent; import javax.annotation.PostConstruct; @@ -122,6 +124,7 @@ public class DefaultTransportService implements TransportService { private final StatsFactory statsFactory; private final TransportDeviceProfileCache deviceProfileCache; private final TransportTenantProfileCache tenantProfileCache; + private final TbUsageStatsClient apiUsageStatsClient; private final TransportRateLimitService rateLimitService; private final DataDecodingEncodingService dataDecodingEncodingService; private final SchedulerComponent scheduler; @@ -150,7 +153,8 @@ public class DefaultTransportService implements TransportService { StatsFactory statsFactory, TransportDeviceProfileCache deviceProfileCache, TransportTenantProfileCache tenantProfileCache, - TransportRateLimitService rateLimitService, DataDecodingEncodingService dataDecodingEncodingService, SchedulerComponent scheduler) { + TbUsageStatsClient apiUsageStatsClient, TransportRateLimitService rateLimitService, + DataDecodingEncodingService dataDecodingEncodingService, SchedulerComponent scheduler) { this.serviceInfoProvider = serviceInfoProvider; this.queueProvider = queueProvider; this.producerProvider = producerProvider; @@ -158,6 +162,7 @@ public class DefaultTransportService implements TransportService { this.statsFactory = statsFactory; this.deviceProfileCache = deviceProfileCache; this.tenantProfileCache = tenantProfileCache; + this.apiUsageStatsClient = apiUsageStatsClient; this.rateLimitService = rateLimitService; this.dataDecodingEncodingService = dataDecodingEncodingService; this.scheduler = scheduler; @@ -362,7 +367,7 @@ public class DefaultTransportService implements TransportService { reportActivityInternal(sessionInfo); TenantId tenantId = new TenantId(new UUID(sessionInfo.getTenantIdMSB(), sessionInfo.getTenantIdLSB())); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); - MsgPackCallback packCallback = new MsgPackCallback(msg.getTsKvListCount(), callback); + MsgPackCallback packCallback = new MsgPackCallback(msg.getTsKvListCount(), new ApiStatsProxyCallback<>(tenantId, dataPoints, callback)); for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { TbMsgMetaData metaData = new TbMsgMetaData(); metaData.putValue("deviceName", sessionInfo.getDeviceName()); @@ -806,4 +811,31 @@ public class DefaultTransportService implements TransportService { callback.onError(t); } } + + private class ApiStatsProxyCallback implements TransportServiceCallback { + private final TenantId tenantId; + private final int dataPoints; + private final TransportServiceCallback callback; + + public ApiStatsProxyCallback(TenantId tenantId, int dataPoints, TransportServiceCallback callback) { + this.tenantId = tenantId; + this.dataPoints = dataPoints; + this.callback = callback; + } + + @Override + public void onSuccess(T msg) { + try { + apiUsageStatsClient.report(tenantId, ApiUsageRecordKey.MSG_COUNT, 1); + apiUsageStatsClient.report(tenantId, ApiUsageRecordKey.DP_TRANSPORT_COUNT, dataPoints); + } finally { + callback.onSuccess(msg); + } + } + + @Override + public void onError(Throwable e) { + callback.onError(e); + } + } }