Browse Source

Api Usage Stats flow

pull/3612/head
Andrii Shvaika 6 years ago
parent
commit
cc2446442d
  1. 105
      application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
  2. 30
      application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
  4. 47
      application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  6. 47
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbUsageStatsClient.java
  7. 36
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

105
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<TenantId, TenantApiUsageState> 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<ToUsageStatsServiceMsg> 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<TsKvEntry> 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<TsKvEntry> 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;
}
}

30
application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbUsageStatsService.java

@ -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<TransportProtos.ToUsageStatsServiceMsg> msg, TbCallback callback) {
}
}

2
application/src/main/java/org/thingsboard/server/service/apiusage/TbUsageStatsService.java → 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<TransportProtos.ToUsageStatsServiceMsg> msg, TbCallback callback);
}

47
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<ApiUsageRecordKey, Long> 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;
}
}

6
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<ToCore
private final TbQueueConsumer<TbProtoQueueMsg<ToCoreMsg>> 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<ToCore
DeviceStateService stateService, TbLocalSubscriptionService localSubscriptionService,
SubscriptionManagerService subscriptionManagerService, DataDecodingEncodingService encodingService,
TbCoreDeviceRpcService tbCoreDeviceRpcService, StatsFactory statsFactory, TbDeviceProfileCache deviceProfileCache,
TbUsageStatsService statsService) {
TbApiUsageStateService statsService) {
super(actorContext, encodingService, deviceProfileCache, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer());
this.mainConsumer = tbCoreQueueFactory.createToCoreMsgConsumer();
this.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer();

47
common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbUsageStatsClient.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.queue.usagestats;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
@ -38,11 +39,12 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
@Component
@Slf4j
public class DefaultTbUsageStatsClient implements TbUsageStatsClient {
@Value("${usage.stats.report.enabled:true}")
private boolean enabled;
@Value("${usage.stats.report.interval:600}")
@Value("${usage.stats.report.interval:10}")
private int interval;
private final ConcurrentMap<TenantId, AtomicLong>[] values = new ConcurrentMap[ApiUsageRecordKey.values().length];
@ -69,28 +71,33 @@ public class DefaultTbUsageStatsClient implements TbUsageStatsClient {
}
private void reportStats() {
ConcurrentMap<TenantId, ToUsageStatsServiceMsg.Builder> report = new ConcurrentHashMap<>();
try {
ConcurrentMap<TenantId, ToUsageStatsServiceMsg.Builder> 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

36
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<T> implements TransportServiceCallback<T> {
private final TenantId tenantId;
private final int dataPoints;
private final TransportServiceCallback<T> callback;
public ApiStatsProxyCallback(TenantId tenantId, int dataPoints, TransportServiceCallback<T> 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);
}
}
}

Loading…
Cancel
Save