Browse Source

refactoring

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
f1ce07ef3f
  1. 34
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
  3. 5
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  4. 36
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  5. 11
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java

34
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric;
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration;
@ -54,6 +53,7 @@ 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.SingleValueArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry;
import org.thingsboard.server.utils.CalculatedFieldArgumentUtils;
import java.util.Collections;
import java.util.HashMap;
@ -145,14 +145,14 @@ public abstract class AbstractCalculatedFieldProcessingService {
));
}
protected ArgumentEntry resolveArgumentValue(String argName, ListenableFuture<ArgumentEntry> future) {
protected ArgumentEntry resolveArgumentValue(String key, ListenableFuture<ArgumentEntry> future) {
try {
return future.get();
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw new RuntimeException("Failed to fetch " + argName + ": " + cause.getMessage(), cause);
throw new RuntimeException("Failed to fetch " + key + ": " + cause.getMessage(), cause);
} catch (InterruptedException e) {
throw new RuntimeException("Failed to fetch" + argName, e);
throw new RuntimeException("Failed to fetch" + key, e);
}
}
@ -298,30 +298,12 @@ public abstract class AbstractCalculatedFieldProcessingService {
};
}
protected ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String metricName, CalculatedFieldCtx ctx) throws Exception {
var config = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration();
AggMetric metric = config.getMetrics().get(metricName);
protected ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval) {
AggFunction function = metric.getFunction();
AggKeyInput input = (AggKeyInput) metric.getInput();
String argName = input.getKey();
Argument argument = ctx.getArguments().get(argName);
String key = argument.getRefEntityKey().getKey();
long intervalMs = interval.getEndTs() - interval.getStartTs();
BaseReadTsKvQuery query = new BaseReadTsKvQuery(key, interval.getStartTs(), interval.getEndTs(), intervalMs, 1, Aggregation.valueOf(function.name()));
log.trace("[{}][{}] Fetching timeseries for query {}", ctx.getTenantId(), entityId, query);
ListenableFuture<List<TsKvEntry>> tsFuture = timeseriesService.findAll(ctx.getTenantId(), entityId, List.of(query));
ListenableFuture<ArgumentEntry> argumentEntryFut = Futures.transform(tsFuture, timeSeries -> {
log.debug("[{}][{}] Fetched {} timeseries for query {}", ctx.getTenantId(), entityId, timeSeries == null ? 0 : timeSeries.size(), query);
if (timeSeries == null || timeSeries.isEmpty()) {
return new SingleValueArgumentEntry();
}
return ArgumentEntry.createSingleValueArgument(timeSeries.get(0));
}, calculatedFieldCallbackExecutor);
return resolveArgumentValue(argName, argumentEntryFut);
BaseReadTsKvQuery query = new BaseReadTsKvQuery(argKey, interval.getStartTs(), interval.getEndTs(), intervalMs, 1, Aggregation.valueOf(function.name()));
ListenableFuture<ArgumentEntry> argumentEntryFut = fetchTimeSeriesInternal(tenantId, entityId, query, CalculatedFieldArgumentUtils::transformAggMetricArgument);
return resolveArgumentValue(argKey, argumentEntryFut);
}
private ListenableFuture<ArgumentEntry> fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) {

3
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.cf;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -38,7 +39,7 @@ public interface CalculatedFieldProcessingService {
Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments);
ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String argName, CalculatedFieldCtx ctx) throws Exception;
ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval);
void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);

5
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -24,6 +24,7 @@ import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -113,8 +114,8 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
}
@Override
public ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String metricName, CalculatedFieldCtx ctx) throws Exception {
return super.fetchMetricDuringInterval(entityId, interval, metricName, ctx);
public ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval) {
return super.fetchMetricDuringInterval(tenantId, entityId, argKey, metric, interval);
}
@Override

36
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java

@ -120,11 +120,9 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results = new HashMap<>();
List<AggIntervalEntry> expiredIntervals = new ArrayList<>();
for (Map.Entry<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> entry : intervals.entrySet()) {
AggIntervalEntry intervalEntry = entry.getKey();
Map<String, AggIntervalEntryStatus> args = entry.getValue();
processInterval(now, intervalEntry, args, expiredIntervals, results);
}
intervals.forEach((intervalEntry, argIntervalStatuses) -> {
processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results);
});
expiredIntervals.forEach(intervals::remove);
ArrayNode result = toResult(results);
@ -165,7 +163,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
AggIntervalEntry intervalEntry,
Map<String, AggIntervalEntryStatus> args,
List<AggIntervalEntry> expiredIntervals,
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) throws Exception {
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) {
long startTs = intervalEntry.getStartTs();
long endTs = intervalEntry.getEndTs();
@ -179,37 +177,35 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
private void handleExpiredInterval(AggIntervalEntry intervalEntry,
Map<String, AggIntervalEntryStatus> args,
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) throws Exception {
for (Map.Entry<String, AggIntervalEntryStatus> argStatus : args.entrySet()) {
String argName = argStatus.getKey();
AggIntervalEntryStatus argEntryIntervalStatus = argStatus.getValue();
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) {
args.forEach((argName, argEntryIntervalStatus) -> {
if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) {
processMetric(intervalEntry, argName, results);
}
}
});
}
private void handleActiveInterval(AggIntervalEntry intervalEntry,
Map<String, AggIntervalEntryStatus> args,
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) throws Exception {
for (Map.Entry<String, AggIntervalEntryStatus> argStatus : args.entrySet()) {
String argName = argStatus.getKey();
AggIntervalEntryStatus argEntryIntervalStatus = argStatus.getValue();
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) {
args.forEach((argName, argEntryIntervalStatus) -> {
if (argEntryIntervalStatus.shouldRecalculate(checkInterval)) {
processMetric(intervalEntry, argName, results);
ctx.scheduleReevaluation(checkInterval, actorCtx);
}
}
});
}
private void processMetric(AggIntervalEntry intervalEntry,
String argName,
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) throws Exception {
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) {
String metricName = findMetricName(argName);
if (metricName != null) {
ArgumentEntry metric = cfProcessingService.fetchMetricDuringInterval(entityId, intervalEntry, metricName, ctx);
if (!metric.isEmpty()) {
results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metric);
AggMetric metric = metrics.get(metricName);
String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey();
ArgumentEntry metricEntry = cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry);
if (!metricEntry.isEmpty()) {
results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry);
}
}
}

11
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java

@ -67,10 +67,17 @@ public class CalculatedFieldArgumentUtils {
return ArgumentEntry.createTsRollingArgument(tsRolling, limit, argTimeWindow);
}
public static ArgumentEntry transformAggregationArgument(List<TsKvEntry> telemetry, long startIntervalTs, long endIntervalTs) {
public static ArgumentEntry transformAggMetricArgument(List<TsKvEntry> timeSeries) {
if (timeSeries == null || timeSeries.isEmpty()) {
return new SingleValueArgumentEntry();
}
return ArgumentEntry.createSingleValueArgument(timeSeries.get(0));
}
public static ArgumentEntry transformAggregationArgument(List<TsKvEntry> timeSeries, long startIntervalTs, long endIntervalTs) {
Map<AggIntervalEntry, AggIntervalEntryStatus> aggIntervals = new HashMap<>();
AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(startIntervalTs, endIntervalTs);
if (telemetry == null || telemetry.isEmpty()) {
if (timeSeries == null || timeSeries.isEmpty()) {
aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus());
} else {
aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus(System.currentTimeMillis()));

Loading…
Cancel
Save