From f9aa9b6a92b0522517ecc7b898951e2e7599094a Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 14 Jul 2023 14:43:32 +0200 Subject: [PATCH] added ability to rewrite latest --- .../controller/TelemetryController.java | 10 +++--- .../DefaultTelemetrySubscriptionService.java | 4 +-- .../dao/timeseries/TimeseriesService.java | 2 ++ .../dao/timeseries/BaseTimeseriesService.java | 33 ++++++++++++++++--- .../api/RuleEngineTelemetryService.java | 2 +- ui-ngx/src/app/core/http/attribute.service.ts | 2 +- 6 files changed, 41 insertions(+), 12 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index 8d94fb22cf..449e82fc41 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/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 deleteLatestTimeseries(EntityId entityIdStr, String keysStr) throws ThingsboardException { + private DeferredResult deleteLatestTimeseries(EntityId entityIdStr, String keysStr, boolean rewrite) throws ThingsboardException { List 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); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 40f4e3415b..a97b7f386d 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/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 keys, FutureCallback callback) { - ListenableFuture> deleteFuture = tsService.removeLatest(tenantId, entityId, keys); + public void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List keys, boolean rewrite, FutureCallback callback) { + ListenableFuture> deleteFuture = tsService.removeLatest(tenantId, entityId, keys, rewrite); addVoidCallback(deleteFuture, callback); addWsCallback(deleteFuture, list -> onTimeSeriesDelete(tenantId, entityId, keys, list)); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index c2bc997235..06e42e09e7 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -58,6 +58,8 @@ public interface TimeseriesService { ListenableFuture> removeLatest(TenantId tenantId, EntityId entityId, Collection keys); + ListenableFuture> removeLatest(TenantId tenantId, EntityId entityId, Collection keys, boolean rewrite); + ListenableFuture> removeAllLatest(TenantId tenantId, EntityId entityId); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 6b8bfa9d64..d101e63a65 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/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> removeLatest(TenantId tenantId, EntityId entityId, Collection keys) { + return removeLatest(tenantId, entityId, keys, false); + } + + @Override + public ListenableFuture> removeLatest(TenantId tenantId, EntityId entityId, Collection keys, boolean rewrite) { validate(entityId); List> 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> latestFuture; + + if (rewrite) { + latestFuture = findLatest(tenantId, entityId, keys); + } else { + latestFuture = Futures.immediateFuture(null); } - return Futures.allAsList(futures); + + return Futures.transformAsync(latestFuture, latest -> { + Map 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 diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java index 9acd03f665..795eaeb785 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java +++ b/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 keys, List deleteTsKvQueries, FutureCallback callback); - void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List keys, FutureCallback callback); + void deleteLatestAndNotify(TenantId tenantId, EntityId entityId, List keys, boolean rewrite, FutureCallback callback); } diff --git a/ui-ngx/src/app/core/http/attribute.service.ts b/ui-ngx/src/app/core/http/attribute.service.ts index cc20069e04..f568758f41 100644 --- a/ui-ngx/src/app/core/http/attribute.service.ts +++ b/ui-ngx/src/app/core/http/attribute.service.ts @@ -68,7 +68,7 @@ export class AttributeService { config?: RequestConfig): Observable { 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)); }