@ -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 ;
@ -138,7 +137,7 @@ public class BaseTimeseriesService implements TimeseriesService {
@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 = new ArrayList ( keys . size ( ) ) ;
List < TsKvEntry > latestEntries = new ArrayList < > ( keys . size ( ) ) ;
keys . forEach ( key - > Validator . validateString ( key , k - > "Incorrect key " + k ) ) ;
keys . forEach ( key - > Validator . validateString ( key , k - > "Incorrect key " + k ) ) ;
for ( String key : keys ) {
for ( String key : keys ) {
latestEntries . add ( timeseriesLatestDao . findLatestSync ( tenantId , entityId , key ) ) ;
latestEntries . add ( timeseriesLatestDao . findLatestSync ( tenantId , entityId , key ) ) ;
@ -170,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 ( ) ) ;
}
}
@ -187,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 ) ;
@ -200,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 ) ) ;
}
}
@ -247,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 ) ;
}
}
@ -257,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 ) ) ;