Browse Source

added ability to rewrite latest

pull/9020/head
YevhenBondarenko 3 years ago
parent
commit
f9aa9b6a92
  1. 10
      application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  3. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  4. 33
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  5. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java
  6. 2
      ui-ngx/src/app/core/http/attribute.service.ts

10
application/src/main/java/org/thingsboard/server/controller/TelemetryController.java

@ -557,12 +557,14 @@ public class TelemetryController extends BaseController {
@ApiParam(value = ENTITY_ID_PARAM_DESCRIPTION, required = true)
@PathVariable("entityId") String entityIdStr,
@ApiParam(value = TELEMETRY_KEYS_DESCRIPTION, required = true)
@RequestParam(name = "keys") String keysStr) throws ThingsboardException {
@RequestParam(name = "keys") String keysStr,
@ApiParam(value = "If the parameter is set to true, the latest telemetry will be rewritten in case that current latest value was removed, otherwise, in case that parameter is set to false the new latest value will not set.")
@RequestParam(name = "rewrite", defaultValue = "false") boolean rewrite) throws ThingsboardException {
EntityId entityId = EntityIdFactory.getByTypeAndId(entityType, entityIdStr);
return deleteLatestTimeseries(entityId, keysStr);
return deleteLatestTimeseries(entityId, keysStr, rewrite);
}
private DeferredResult<ResponseEntity> deleteLatestTimeseries(EntityId entityIdStr, String keysStr) throws ThingsboardException {
private DeferredResult<ResponseEntity> deleteLatestTimeseries(EntityId entityIdStr, String keysStr, boolean rewrite) throws ThingsboardException {
List<String> keys = toKeysList(keysStr);
if (keys.isEmpty()) {
return getImmediateDeferredResult("Empty keys: " + keysStr, HttpStatus.BAD_REQUEST);
@ -570,7 +572,7 @@ public class TelemetryController extends BaseController {
SecurityUser user = getCurrentUser();
return accessValidator.validateEntityAndCallback(user, Operation.WRITE_TELEMETRY, entityIdStr, (result, tenantId, entityId) ->
tsSubService.deleteLatestAndNotify(tenantId, entityId, keys, new FutureCallback<>() {
tsSubService.deleteLatestAndNotify(tenantId, entityId, keys, rewrite, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable Void tmp) {
logLatestTimeseriesDeleted(user, entityId, keys, null);

4
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -317,8 +317,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}
@Override
public void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback) {
ListenableFuture<List<TsKvLatestRemovingResult>> deleteFuture = tsService.removeLatest(tenantId, entityId, keys);
public void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List<String> keys, boolean rewrite, FutureCallback<Void> callback) {
ListenableFuture<List<TsKvLatestRemovingResult>> deleteFuture = tsService.removeLatest(tenantId, entityId, keys, rewrite);
addVoidCallback(deleteFuture, callback);
addWsCallback(deleteFuture, list -> onTimeSeriesDelete(tenantId, entityId, keys, list));
}

2
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -58,6 +58,8 @@ public interface TimeseriesService {
ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys);
ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys, boolean rewrite);
ListenableFuture<Collection<String>> removeAllLatest(TenantId tenantId, EntityId entityId);
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);

33
dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java

@ -46,6 +46,7 @@ import org.thingsboard.server.dao.service.Validator;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
@ -251,13 +252,37 @@ public class BaseTimeseriesService implements TimeseriesService {
@Override
public ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) {
return removeLatest(tenantId, entityId, keys, false);
}
@Override
public ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys, boolean rewrite) {
validate(entityId);
List<ListenableFuture<TsKvLatestRemovingResult>> futures = Lists.newArrayListWithExpectedSize(keys.size());
for (String key : keys) {
DeleteTsKvQuery query = new BaseDeleteTsKvQuery(key, 0, System.currentTimeMillis(), false);
futures.add(timeseriesLatestDao.removeLatest(tenantId, entityId, query));
ListenableFuture<List<TsKvEntry>> latestFuture;
if (rewrite) {
latestFuture = findLatest(tenantId, entityId, keys);
} else {
latestFuture = Futures.immediateFuture(null);
}
return Futures.allAsList(futures);
return Futures.transformAsync(latestFuture, latest -> {
Map<String, Long> keyTsMap;
if (latest != null) {
keyTsMap = latest.stream().collect(Collectors.toMap(TsKvEntry::getKey, TsKvEntry::getTs));
} else {
keyTsMap = Collections.emptyMap();
}
for (String key : keys) {
long startTs = keyTsMap.getOrDefault(key, 0L);
DeleteTsKvQuery query = new BaseDeleteTsKvQuery(key, startTs, System.currentTimeMillis(), rewrite);
futures.add(timeseriesLatestDao.removeLatest(tenantId, entityId, query));
}
return Futures.allAsList(futures);
}, MoreExecutors.directExecutor());
}
@Override

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java

@ -72,5 +72,5 @@ public interface RuleEngineTelemetryService {
void deleteTimeseriesAndNotify(TenantId tenantId, EntityId entityId, List<String> keys, List<DeleteTsKvQuery> deleteTsKvQueries, FutureCallback<Void> callback);
void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback);
void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List<String> keys, boolean rewrite, FutureCallback<Void> callback);
}

2
ui-ngx/src/app/core/http/attribute.service.ts

@ -68,7 +68,7 @@ export class AttributeService {
config?: RequestConfig): Observable<any> {
const keys = timeseries.map(attribute => encodeURIComponent(attribute.key)).join(',');
let url = `/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/timeseries/latest/delete?keys=${keys}` +
`$rewrite=${rewrite}`;
`&rewrite=${rewrite}`;
return this.http.delete(url, defaultHttpOptionsFromConfig(config));
}

Loading…
Cancel
Save