|
|
@ -16,7 +16,6 @@ |
|
|
package org.thingsboard.server.dao.timeseries; |
|
|
package org.thingsboard.server.dao.timeseries; |
|
|
|
|
|
|
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.collect.Lists; |
|
|
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.MoreExecutors; |
|
|
import com.google.common.util.concurrent.MoreExecutors; |
|
|
@ -43,6 +42,7 @@ import org.thingsboard.server.dao.entityview.EntityViewService; |
|
|
import org.thingsboard.server.dao.exception.IncorrectParameterException; |
|
|
import org.thingsboard.server.dao.exception.IncorrectParameterException; |
|
|
import org.thingsboard.server.dao.service.Validator; |
|
|
import org.thingsboard.server.dao.service.Validator; |
|
|
|
|
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
import java.util.Collection; |
|
|
import java.util.Collection; |
|
|
import java.util.Collections; |
|
|
import java.util.Collections; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
@ -126,18 +126,22 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
public ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
validate(entityId); |
|
|
validate(entityId); |
|
|
List<ListenableFuture<TsKvEntry>> futures = Lists.newArrayListWithExpectedSize(keys.size()); |
|
|
List<ListenableFuture<TsKvEntry>> futures = new ArrayList<>(keys.size()); |
|
|
keys.forEach(key -> Validator.validateString(key, "Incorrect key " + key)); |
|
|
keys.forEach(key -> Validator.validateString(key, k -> "Incorrect key " + k)); |
|
|
keys.forEach(key -> futures.add(timeseriesLatestDao.findLatest(tenantId, entityId, key))); |
|
|
for (String key : keys) { |
|
|
|
|
|
futures.add(timeseriesLatestDao.findLatest(tenantId, entityId, key)); |
|
|
|
|
|
} |
|
|
return Futures.allAsList(futures); |
|
|
return Futures.allAsList(futures); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public List<TsKvEntry> findLatestSync(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
public List<TsKvEntry> findLatestSync(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
validate(entityId); |
|
|
validate(entityId); |
|
|
List<TsKvEntry> latestEntries = Lists.newArrayListWithExpectedSize(keys.size()); |
|
|
List<TsKvEntry> latestEntries = new ArrayList<>(keys.size()); |
|
|
keys.forEach(key -> Validator.validateString(key, "Incorrect key " + key)); |
|
|
keys.forEach(key -> Validator.validateString(key, k -> "Incorrect key " + k)); |
|
|
keys.forEach(key -> latestEntries.add(timeseriesLatestDao.findLatestSync(tenantId, entityId, key))); |
|
|
for (String key : keys) { |
|
|
|
|
|
latestEntries.add(timeseriesLatestDao.findLatestSync(tenantId, entityId, key)); |
|
|
|
|
|
} |
|
|
return latestEntries; |
|
|
return latestEntries; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -165,7 +169,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { |
|
|
public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { |
|
|
validate(entityId); |
|
|
validate(entityId); |
|
|
List<ListenableFuture<Integer>> futures = Lists.newArrayListWithExpectedSize(INSERTS_PER_ENTRY); |
|
|
List<ListenableFuture<Integer>> futures = new ArrayList<>(INSERTS_PER_ENTRY); |
|
|
saveAndRegisterFutures(tenantId, futures, entityId, tsKvEntry, 0L); |
|
|
saveAndRegisterFutures(tenantId, futures, entityId, tsKvEntry, 0L); |
|
|
return Futures.transform(Futures.allAsList(futures), SUM_ALL_INTEGERS, MoreExecutors.directExecutor()); |
|
|
return Futures.transform(Futures.allAsList(futures), SUM_ALL_INTEGERS, MoreExecutors.directExecutor()); |
|
|
} |
|
|
} |
|
|
@ -182,7 +186,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
|
|
|
|
|
|
private ListenableFuture<Integer> doSave(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries, long ttl, boolean saveLatest) { |
|
|
private ListenableFuture<Integer> doSave(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries, long ttl, boolean saveLatest) { |
|
|
int inserts = saveLatest ? INSERTS_PER_ENTRY : INSERTS_PER_ENTRY_WITHOUT_LATEST; |
|
|
int inserts = saveLatest ? INSERTS_PER_ENTRY : INSERTS_PER_ENTRY_WITHOUT_LATEST; |
|
|
List<ListenableFuture<Integer>> futures = Lists.newArrayListWithExpectedSize(tsKvEntries.size() * inserts); |
|
|
List<ListenableFuture<Integer>> futures = new ArrayList<>(tsKvEntries.size() * inserts); |
|
|
for (TsKvEntry tsKvEntry : tsKvEntries) { |
|
|
for (TsKvEntry tsKvEntry : tsKvEntries) { |
|
|
if (saveLatest) { |
|
|
if (saveLatest) { |
|
|
saveAndRegisterFutures(tenantId, futures, entityId, tsKvEntry, ttl); |
|
|
saveAndRegisterFutures(tenantId, futures, entityId, tsKvEntry, ttl); |
|
|
@ -195,7 +199,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<List<Void>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries) { |
|
|
public ListenableFuture<List<Void>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries) { |
|
|
List<ListenableFuture<Void>> futures = Lists.newArrayListWithExpectedSize(tsKvEntries.size()); |
|
|
List<ListenableFuture<Void>> futures = new ArrayList<>(tsKvEntries.size()); |
|
|
for (TsKvEntry tsKvEntry : tsKvEntries) { |
|
|
for (TsKvEntry tsKvEntry : tsKvEntries) { |
|
|
futures.add(timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry)); |
|
|
futures.add(timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry)); |
|
|
} |
|
|
} |
|
|
@ -242,7 +246,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
public ListenableFuture<List<TsKvLatestRemovingResult>> remove(TenantId tenantId, EntityId entityId, List<DeleteTsKvQuery> deleteTsKvQueries) { |
|
|
public ListenableFuture<List<TsKvLatestRemovingResult>> remove(TenantId tenantId, EntityId entityId, List<DeleteTsKvQuery> deleteTsKvQueries) { |
|
|
validate(entityId); |
|
|
validate(entityId); |
|
|
deleteTsKvQueries.forEach(BaseTimeseriesService::validate); |
|
|
deleteTsKvQueries.forEach(BaseTimeseriesService::validate); |
|
|
List<ListenableFuture<TsKvLatestRemovingResult>> futures = Lists.newArrayListWithExpectedSize(deleteTsKvQueries.size() * DELETES_PER_ENTRY); |
|
|
List<ListenableFuture<TsKvLatestRemovingResult>> futures = new ArrayList<>(deleteTsKvQueries.size() * DELETES_PER_ENTRY); |
|
|
for (DeleteTsKvQuery tsKvQuery : deleteTsKvQueries) { |
|
|
for (DeleteTsKvQuery tsKvQuery : deleteTsKvQueries) { |
|
|
deleteAndRegisterFutures(tenantId, futures, entityId, tsKvQuery); |
|
|
deleteAndRegisterFutures(tenantId, futures, entityId, tsKvQuery); |
|
|
} |
|
|
} |
|
|
@ -252,7 +256,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
public ListenableFuture<List<TsKvLatestRemovingResult>> removeLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) { |
|
|
validate(entityId); |
|
|
validate(entityId); |
|
|
List<ListenableFuture<TsKvLatestRemovingResult>> futures = Lists.newArrayListWithExpectedSize(keys.size()); |
|
|
List<ListenableFuture<TsKvLatestRemovingResult>> futures = new ArrayList<>(keys.size()); |
|
|
for (String key : keys) { |
|
|
for (String key : keys) { |
|
|
DeleteTsKvQuery query = new BaseDeleteTsKvQuery(key, 0, System.currentTimeMillis(), false); |
|
|
DeleteTsKvQuery query = new BaseDeleteTsKvQuery(key, 0, System.currentTimeMillis(), false); |
|
|
futures.add(timeseriesLatestDao.removeLatest(tenantId, entityId, query)); |
|
|
futures.add(timeseriesLatestDao.removeLatest(tenantId, entityId, query)); |
|
|
@ -281,7 +285,7 @@ public class BaseTimeseriesService implements TimeseriesService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private static void validate(EntityId entityId) { |
|
|
private static void validate(EntityId entityId) { |
|
|
Validator.validateEntityId(entityId, "Incorrect entityId " + entityId); |
|
|
Validator.validateEntityId(entityId, id -> "Incorrect entityId " + id); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void validate(ReadTsKvQuery query) { |
|
|
private void validate(ReadTsKvQuery query) { |
|
|
|