|
|
|
@ -19,6 +19,7 @@ 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.ListenableFuture; |
|
|
|
import com.google.common.util.concurrent.SettableFuture; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.springframework.data.domain.PageRequest; |
|
|
|
@ -39,6 +40,8 @@ import org.thingsboard.server.dao.util.SqlDao; |
|
|
|
import javax.annotation.Nullable; |
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.concurrent.CompletableFuture; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID; |
|
|
|
@ -80,7 +83,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
return findAllAsyncWithLimit(entityId, query); |
|
|
|
} else { |
|
|
|
long stepTs = query.getStartTs(); |
|
|
|
List<ListenableFuture<TsKvEntry>> futures = new ArrayList<>(); |
|
|
|
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>(); |
|
|
|
while (stepTs < query.getEndTs()) { |
|
|
|
long startTs = stepTs; |
|
|
|
long endTs = stepTs + query.getInterval(); |
|
|
|
@ -88,16 +91,30 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
futures.add(findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation())); |
|
|
|
stepTs = endTs; |
|
|
|
} |
|
|
|
return Futures.allAsList(futures); |
|
|
|
ListenableFuture<List<Optional<TsKvEntry>>> future = Futures.allAsList(futures); |
|
|
|
return Futures.transform(future, new Function<List<Optional<TsKvEntry>>, List<TsKvEntry>>() { |
|
|
|
@Nullable |
|
|
|
@Override |
|
|
|
public List<TsKvEntry> apply(@Nullable List<Optional<TsKvEntry>> results) { |
|
|
|
if (results == null || results.isEmpty()) { |
|
|
|
return null; |
|
|
|
} |
|
|
|
return results.stream() |
|
|
|
.filter(Optional::isPresent) |
|
|
|
.map(Optional::get) |
|
|
|
.collect(Collectors.toList()); |
|
|
|
} |
|
|
|
}, service); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<TsKvEntry> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { |
|
|
|
TsKvEntity entity; |
|
|
|
private ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { |
|
|
|
CompletableFuture<TsKvEntity> entity; |
|
|
|
String entityIdStr = fromTimeUUID(entityId.getId()); |
|
|
|
switch (aggregation) { |
|
|
|
case AVG: |
|
|
|
entity = tsKvRepository.findAvg( |
|
|
|
fromTimeUUID(entityId.getId()), |
|
|
|
entityIdStr, |
|
|
|
entityId.getEntityType(), |
|
|
|
key, |
|
|
|
startTs, |
|
|
|
@ -106,7 +123,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
break; |
|
|
|
case MAX: |
|
|
|
entity = tsKvRepository.findMax( |
|
|
|
fromTimeUUID(entityId.getId()), |
|
|
|
entityIdStr, |
|
|
|
entityId.getEntityType(), |
|
|
|
key, |
|
|
|
startTs, |
|
|
|
@ -115,7 +132,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
break; |
|
|
|
case MIN: |
|
|
|
entity = tsKvRepository.findMin( |
|
|
|
fromTimeUUID(entityId.getId()), |
|
|
|
entityIdStr, |
|
|
|
entityId.getEntityType(), |
|
|
|
key, |
|
|
|
startTs, |
|
|
|
@ -124,7 +141,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
break; |
|
|
|
case SUM: |
|
|
|
entity = tsKvRepository.findSum( |
|
|
|
fromTimeUUID(entityId.getId()), |
|
|
|
entityIdStr, |
|
|
|
entityId.getEntityType(), |
|
|
|
key, |
|
|
|
startTs, |
|
|
|
@ -133,7 +150,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
break; |
|
|
|
case COUNT: |
|
|
|
entity = tsKvRepository.findCount( |
|
|
|
fromTimeUUID(entityId.getId()), |
|
|
|
entityIdStr, |
|
|
|
entityId.getEntityType(), |
|
|
|
key, |
|
|
|
startTs, |
|
|
|
@ -141,12 +158,32 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp |
|
|
|
|
|
|
|
break; |
|
|
|
default: |
|
|
|
entity = null; |
|
|
|
throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); |
|
|
|
} |
|
|
|
if (entity != null) { |
|
|
|
entity.setTs(ts); |
|
|
|
} |
|
|
|
return service.submit(() -> DaoUtil.getData(entity)); |
|
|
|
|
|
|
|
SettableFuture<TsKvEntity> listenableFuture = SettableFuture.create(); |
|
|
|
entity.whenComplete((tsKvEntity, throwable) -> { |
|
|
|
if (throwable != null) { |
|
|
|
listenableFuture.setException(throwable); |
|
|
|
} else { |
|
|
|
listenableFuture.set(tsKvEntity); |
|
|
|
} |
|
|
|
}); |
|
|
|
return Futures.transform(listenableFuture, new Function<TsKvEntity, Optional<TsKvEntry>>() { |
|
|
|
@Nullable |
|
|
|
@Override |
|
|
|
public Optional<TsKvEntry> apply(@Nullable TsKvEntity entity) { |
|
|
|
if (entity != null && entity.isNotEmpty()) { |
|
|
|
entity.setEntityId(entityIdStr); |
|
|
|
entity.setEntityType(entityId.getEntityType()); |
|
|
|
entity.setKey(key); |
|
|
|
entity.setTs(ts); |
|
|
|
return Optional.of(DaoUtil.getData(entity)); |
|
|
|
} else { |
|
|
|
return Optional.empty(); |
|
|
|
} |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, TsKvQuery query) { |
|
|
|
|