|
|
@ -54,8 +54,6 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntryStatus; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.EntityAggregationArgumentEntry; |
|
|
|
|
|
|
|
|
|
|
|
import java.util.Collections; |
|
|
import java.util.Collections; |
|
|
import java.util.HashMap; |
|
|
import java.util.HashMap; |
|
|
@ -64,7 +62,7 @@ import java.util.Map; |
|
|
import java.util.Optional; |
|
|
import java.util.Optional; |
|
|
import java.util.Set; |
|
|
import java.util.Set; |
|
|
import java.util.concurrent.ExecutionException; |
|
|
import java.util.concurrent.ExecutionException; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
import java.util.function.Function; |
|
|
import java.util.stream.Collectors; |
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
|
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION; |
|
|
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION; |
|
|
@ -73,7 +71,9 @@ import static org.thingsboard.server.common.data.cf.configuration.geofencing.Ent |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
|
|
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformAggregationArgument; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
|
|
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformTsRollingArgument; |
|
|
|
|
|
|
|
|
@Data |
|
|
@Data |
|
|
@Slf4j |
|
|
@Slf4j |
|
|
@ -141,19 +141,21 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
return argFutures.entrySet().stream() |
|
|
return argFutures.entrySet().stream() |
|
|
.collect(Collectors.toMap( |
|
|
.collect(Collectors.toMap( |
|
|
Map.Entry::getKey, // Keep the key as is
|
|
|
Map.Entry::getKey, // Keep the key as is
|
|
|
entry -> { |
|
|
entry -> resolveArgumentValue(entry.getKey(), entry.getValue()) |
|
|
try { |
|
|
|
|
|
return entry.getValue().get(); |
|
|
|
|
|
} catch (ExecutionException e) { |
|
|
|
|
|
Throwable cause = e.getCause(); |
|
|
|
|
|
throw new RuntimeException("Failed to fetch " + entry.getKey() + ": " + cause.getMessage(), cause); |
|
|
|
|
|
} catch (InterruptedException e) { |
|
|
|
|
|
throw new RuntimeException("Failed to fetch" + entry.getKey(), e); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
)); |
|
|
)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
protected ArgumentEntry resolveArgumentValue(String argName, ListenableFuture<ArgumentEntry> future) { |
|
|
|
|
|
try { |
|
|
|
|
|
return future.get(); |
|
|
|
|
|
} catch (ExecutionException e) { |
|
|
|
|
|
Throwable cause = e.getCause(); |
|
|
|
|
|
throw new RuntimeException("Failed to fetch " + argName + ": " + cause.getMessage(), cause); |
|
|
|
|
|
} catch (InterruptedException e) { |
|
|
|
|
|
throw new RuntimeException("Failed to fetch" + argName, e); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
protected ListenableFuture<ArgumentEntry> fetchPropagationCalculatedFieldArgument(CalculatedFieldCtx ctx, EntityId entityId) { |
|
|
protected ListenableFuture<ArgumentEntry> fetchPropagationCalculatedFieldArgument(CalculatedFieldCtx ctx, EntityId entityId) { |
|
|
ListenableFuture<List<EntityId>> propagationEntityIds = fromDynamicSource(ctx.getTenantId(), entityId, ctx.getPropagationArgument()); |
|
|
ListenableFuture<List<EntityId>> propagationEntityIds = fromDynamicSource(ctx.getTenantId(), entityId, ctx.getPropagationArgument()); |
|
|
return Futures.transform(propagationEntityIds, ArgumentEntry::createPropagationArgument, MoreExecutors.directExecutor()); |
|
|
return Futures.transform(propagationEntityIds, ArgumentEntry::createPropagationArgument, MoreExecutors.directExecutor()); |
|
|
@ -199,7 +201,7 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
return aggConfig.getArguments().entrySet().stream() |
|
|
return aggConfig.getArguments().entrySet().stream() |
|
|
.collect(Collectors.toMap( |
|
|
.collect(Collectors.toMap( |
|
|
Map.Entry::getKey, |
|
|
Map.Entry::getKey, |
|
|
entry -> fetchTimeSeries(ctx.getTenantId(), entityId, entry.getValue(), aggConfig.getInterval()) |
|
|
entry -> fetchTimeSeries(ctx.getTenantId(), entityId, entry.getValue(), aggConfig.getInterval(), ts) |
|
|
)); |
|
|
)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -319,45 +321,21 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
return ArgumentEntry.createSingleValueArgument(timeSeries.get(0)); |
|
|
return ArgumentEntry.createSingleValueArgument(timeSeries.get(0)); |
|
|
}, calculatedFieldCallbackExecutor); |
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
|
|
|
|
|
// Ugly but necessary. We do not expect to often fetch data from DB. Only once per <Entity, CalculatedField> pair lifetime.
|
|
|
return resolveArgumentValue(argName, argumentEntryFut); |
|
|
// This call happens while processing the CF pack from the queue consumer. So the timeout should be relatively low.
|
|
|
|
|
|
// Alternatively, we can fetch the state outside the actor system and push separate command to create this actor,
|
|
|
|
|
|
// but this will significantly complicate the code.
|
|
|
|
|
|
return argumentEntryFut.get(1, TimeUnit.MINUTES); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval) { |
|
|
private ListenableFuture<ArgumentEntry> fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { |
|
|
long startInterval = interval.getCurrentIntervalStartTs(); |
|
|
long startInterval = interval.getCurrentIntervalStartTs(); |
|
|
|
|
|
long intervalEndTs = interval.getCurrentIntervalEndTs(); |
|
|
String key = argument.getRefEntityKey().getKey(); |
|
|
ReadTsKvQuery query = buildTimeSeriesQuery(tenantId, argument, startInterval, queryEndTs); |
|
|
ReadTsKvQuery query = new BaseReadTsKvQuery(key, startInterval, System.currentTimeMillis(), 0, 1, Aggregation.NONE); |
|
|
return fetchTimeSeriesInternal(tenantId, entityId, query, timeSeries -> transformAggregationArgument(timeSeries, startInterval, intervalEndTs)); |
|
|
|
|
|
|
|
|
log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query); |
|
|
|
|
|
ListenableFuture<List<TsKvEntry>> fetchedTelemetryFut = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
|
|
|
|
|
return Futures.transform(fetchedTelemetryFut, telemetry -> { |
|
|
|
|
|
log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, telemetry == null ? 0 : telemetry.size(), query); |
|
|
|
|
|
Map<AggIntervalEntry, AggIntervalEntryStatus> aggIntervals = new HashMap<>(); |
|
|
|
|
|
AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(interval.getCurrentIntervalStartTs(), interval.getCurrentIntervalEndTs()); |
|
|
|
|
|
if (telemetry == null || telemetry.isEmpty()) { |
|
|
|
|
|
aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus()); |
|
|
|
|
|
} else { |
|
|
|
|
|
aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus(System.currentTimeMillis())); |
|
|
|
|
|
} |
|
|
|
|
|
return new EntityAggregationArgumentEntry(aggIntervals); |
|
|
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { |
|
|
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { |
|
|
long argTimeWindow = argument.getTimeWindow() == 0 ? queryEndTs : argument.getTimeWindow(); |
|
|
long argTimeWindow = argument.getTimeWindow() == 0 ? queryEndTs : argument.getTimeWindow(); |
|
|
long startInterval = queryEndTs - argTimeWindow; |
|
|
long startInterval = queryEndTs - argTimeWindow; |
|
|
ReadTsKvQuery query = buildTsRollingQuery(tenantId, argument, startInterval, queryEndTs); |
|
|
ReadTsKvQuery query = buildTimeSeriesQuery(tenantId, argument, startInterval, queryEndTs); |
|
|
|
|
|
return fetchTimeSeriesInternal(tenantId, entityId, query, tsRolling -> transformTsRollingArgument(tsRolling, query.getLimit(), argTimeWindow)); |
|
|
log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query); |
|
|
|
|
|
ListenableFuture<List<TsKvEntry>> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
|
|
|
|
|
return Futures.transform(tsRollingFuture, tsRolling -> { |
|
|
|
|
|
log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, tsRolling == null ? 0 : tsRolling.size(), query); |
|
|
|
|
|
return ArgumentEntry.createTsRollingArgument(tsRolling, query.getLimit(), argTimeWindow); |
|
|
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchAttribute(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) { |
|
|
private ListenableFuture<ArgumentEntry> fetchAttribute(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) { |
|
|
@ -383,7 +361,16 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
}, calculatedFieldCallbackExecutor)); |
|
|
}, calculatedFieldCallbackExecutor)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ReadTsKvQuery buildTsRollingQuery(TenantId tenantId, Argument argument, long startTs, long endTs) { |
|
|
private ListenableFuture<ArgumentEntry> fetchTimeSeriesInternal(TenantId tenantId, EntityId entityId, ReadTsKvQuery query, Function<List<TsKvEntry>, ArgumentEntry> transformArgument) { |
|
|
|
|
|
log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query); |
|
|
|
|
|
ListenableFuture<List<TsKvEntry>> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
|
|
|
|
|
return Futures.transform(tsRollingFuture, tsRolling -> { |
|
|
|
|
|
log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, tsRolling == null ? 0 : tsRolling.size(), query); |
|
|
|
|
|
return transformArgument.apply(tsRolling); |
|
|
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private ReadTsKvQuery buildTimeSeriesQuery(TenantId tenantId, Argument argument, long startTs, long endTs) { |
|
|
long maxDataPoints = apiLimitService.getLimit( |
|
|
long maxDataPoints = apiLimitService.getLimit( |
|
|
tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); |
|
|
tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); |
|
|
int argumentLimit = argument.getLimit(); |
|
|
int argumentLimit = argument.getLimit(); |
|
|
|