|
|
|
@ -33,10 +33,12 @@ import org.thingsboard.rule.engine.api.AttributesSaveRequest; |
|
|
|
import org.thingsboard.rule.engine.api.AttributesSaveRequest.Strategy; |
|
|
|
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; |
|
|
|
import org.thingsboard.server.common.adaptor.JsonConverter; |
|
|
|
import org.thingsboard.server.common.data.AttributeScope; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
|
|
|
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.AttributeImmediateOutputStrategy; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.TimeSeriesImmediateOutputStrategy; |
|
|
|
@ -357,8 +359,8 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
|
|
|
|
SettableFuture<Void> future = SettableFuture.create(); |
|
|
|
switch (type) { |
|
|
|
case ATTRIBUTES -> saveAttributes(tenantId, entityId, cfResult, cfIds, future); |
|
|
|
case TIME_SERIES -> saveTimeSeries(tenantId, entityId, cfResult, cfIds, System.currentTimeMillis(), future); |
|
|
|
case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfIds, future); |
|
|
|
case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), future); |
|
|
|
} |
|
|
|
|
|
|
|
Futures.addCallback(future, new FutureCallback<>() { |
|
|
|
@ -376,39 +378,37 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
}, MoreExecutors.directExecutor()); |
|
|
|
} |
|
|
|
|
|
|
|
private void saveAttributes(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List<CalculatedFieldId> cfIds, SettableFuture<Void> future) { |
|
|
|
if (!(cfResult.getOutputStrategy() instanceof AttributeImmediateOutputStrategy outputStrategy)) { |
|
|
|
future.setException(new IllegalArgumentException("Expected AttributeImmediateOutputStrategy")); |
|
|
|
return; |
|
|
|
} |
|
|
|
JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); |
|
|
|
private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List<CalculatedFieldId> cfIds, SettableFuture<Void> future) { |
|
|
|
if (!(outputStrategy instanceof AttributeImmediateOutputStrategy attOutputStrategy)) { |
|
|
|
future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); |
|
|
|
} else { |
|
|
|
AttributesSaveRequest.Strategy strategy = new Strategy(attOutputStrategy.isSaveAttribute(), attOutputStrategy.isSendWsUpdate(), attOutputStrategy.isProcessCfs()); |
|
|
|
List<AttributeKvEntry> newAttributes = JsonConverter.convertToAttributes(jsonResult); |
|
|
|
|
|
|
|
AttributesSaveRequest.Strategy strategy = new Strategy(outputStrategy.isSaveAttribute(), outputStrategy.isSendWsUpdate(), outputStrategy.isProcessCfs()); |
|
|
|
List<AttributeKvEntry> newAttributes = JsonConverter.convertToAttributes(jsonResult); |
|
|
|
if (!attOutputStrategy.isUpdateAttributesOnlyOnValueChange()) { |
|
|
|
saveAttributesInternal(tenantId, entityId, scope, cfIds, newAttributes, strategy, future); |
|
|
|
return; |
|
|
|
} |
|
|
|
|
|
|
|
if (!outputStrategy.isUpdateAttributesOnlyOnValueChange()) { |
|
|
|
saveAttributesInternal(tenantId, entityId, cfResult, cfIds, newAttributes, strategy, future); |
|
|
|
return; |
|
|
|
} |
|
|
|
List<String> keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); |
|
|
|
ListenableFuture<List<AttributeKvEntry>> findFuture = attributesService.find(tenantId, entityId, scope, keys); |
|
|
|
|
|
|
|
List<String> keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); |
|
|
|
ListenableFuture<List<AttributeKvEntry>> findFuture = attributesService.find(tenantId, entityId, cfResult.getScope(), keys); |
|
|
|
|
|
|
|
DonAsynchron.withCallback(findFuture, |
|
|
|
existingAttributes -> { |
|
|
|
List<AttributeKvEntry> changed = filterChangedAttr(existingAttributes, newAttributes); |
|
|
|
if (changed.isEmpty()) { |
|
|
|
future.set(null); |
|
|
|
return; |
|
|
|
} |
|
|
|
saveAttributesInternal(tenantId, entityId, cfResult, cfIds, changed, strategy, future); |
|
|
|
}, |
|
|
|
future::setException, |
|
|
|
MoreExecutors.directExecutor()); |
|
|
|
DonAsynchron.withCallback(findFuture, |
|
|
|
existingAttributes -> { |
|
|
|
List<AttributeKvEntry> changed = filterChangedAttr(existingAttributes, newAttributes); |
|
|
|
if (changed.isEmpty()) { |
|
|
|
future.set(null); |
|
|
|
return; |
|
|
|
} |
|
|
|
saveAttributesInternal(tenantId, entityId, scope, cfIds, changed, strategy, future); |
|
|
|
}, |
|
|
|
future::setException, |
|
|
|
MoreExecutors.directExecutor()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void saveAttributesInternal(TenantId tenantId, EntityId entityId, |
|
|
|
TelemetryCalculatedFieldResult cfResult, |
|
|
|
AttributeScope scope, |
|
|
|
List<CalculatedFieldId> cfIds, |
|
|
|
List<AttributeKvEntry> entries, |
|
|
|
AttributesSaveRequest.Strategy strategy, |
|
|
|
@ -416,7 +416,7 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
tsSubService.saveAttributesInternal(AttributesSaveRequest.builder() |
|
|
|
.tenantId(tenantId) |
|
|
|
.entityId(entityId) |
|
|
|
.scope(cfResult.getScope()) |
|
|
|
.scope(scope) |
|
|
|
.entries(entries) |
|
|
|
.strategy(strategy) |
|
|
|
.previousCalculatedFieldIds(cfIds) |
|
|
|
@ -424,28 +424,35 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
.build()); |
|
|
|
} |
|
|
|
|
|
|
|
private void saveTimeSeries(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List<CalculatedFieldId> cfIds, long ts, SettableFuture<Void> future) { |
|
|
|
if (!(cfResult.getOutputStrategy() instanceof TimeSeriesImmediateOutputStrategy outputStrategy)) { |
|
|
|
future.setException(new IllegalArgumentException("Expected TimeSeriesImmediateOutputStrategy")); |
|
|
|
return; |
|
|
|
private void saveTimeSeries(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, List<CalculatedFieldId> cfIds, long ts, SettableFuture<Void> future) { |
|
|
|
if (!(outputStrategy instanceof TimeSeriesImmediateOutputStrategy tsOutputStrategy)) { |
|
|
|
future.setException(new IllegalArgumentException("Only TimeSeriesImmediateOutputStrategy is supported.")); |
|
|
|
} else { |
|
|
|
TimeseriesSaveRequest.Strategy strategy = new TimeseriesSaveRequest.Strategy(tsOutputStrategy.isSaveTimeSeries(), tsOutputStrategy.isSaveLatest(), tsOutputStrategy.isSendWsUpdate(), tsOutputStrategy.isProcessCfs()); |
|
|
|
saveTimeSeriesInternal(tenantId, entityId, jsonResult, tsOutputStrategy.getTtl(), cfIds, ts, strategy, future); |
|
|
|
} |
|
|
|
JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); |
|
|
|
} |
|
|
|
|
|
|
|
private void saveTimeSeriesInternal(TenantId tenantId, EntityId entityId, JsonElement jsonResult, Long ttl, List<CalculatedFieldId> cfIds, long ts, TimeseriesSaveRequest.Strategy strategy, SettableFuture<Void> future) { |
|
|
|
Map<Long, List<KvEntry>> tsKvMap = JsonConverter.convertToTelemetry(jsonResult, ts); |
|
|
|
if (tsKvMap.isEmpty()) { |
|
|
|
future.set(null); |
|
|
|
return; |
|
|
|
} |
|
|
|
List<TsKvEntry> tsEntries = toTsKvEntryList(tsKvMap); |
|
|
|
TimeseriesSaveRequest.Strategy strategy = new TimeseriesSaveRequest.Strategy(outputStrategy.isSaveTimeSeries(), outputStrategy.isSaveLatest(), outputStrategy.isSendWsUpdate(), outputStrategy.isProcessCfs()); |
|
|
|
tsSubService.saveTimeseriesInternal(TimeseriesSaveRequest.builder() |
|
|
|
TimeseriesSaveRequest.Builder builder = TimeseriesSaveRequest.builder() |
|
|
|
.tenantId(tenantId) |
|
|
|
.entityId(entityId) |
|
|
|
.entries(tsEntries) |
|
|
|
.ttl(outputStrategy.getTtl()) |
|
|
|
.strategy(strategy) |
|
|
|
.previousCalculatedFieldIds(cfIds) |
|
|
|
.future(future) |
|
|
|
.build()); |
|
|
|
.future(future); |
|
|
|
if (ttl != null) { |
|
|
|
builder.ttl(ttl); |
|
|
|
} |
|
|
|
if (cfIds != null && !cfIds.isEmpty()) { |
|
|
|
builder.previousCalculatedFieldIds(cfIds); |
|
|
|
} |
|
|
|
tsSubService.saveTimeseriesInternal(builder.build()); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|