Browse Source

changed CF links config

pull/12404/head
IrynaMatveieva 2 years ago
parent
commit
e2ac2708b6
  1. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java
  2. 37
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  3. 28
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  4. 21
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java
  5. 4
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java
  6. 16
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java
  7. 7
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  8. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java
  9. 17
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java
  10. 7
      dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java
  11. 2
      dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java
  12. 3
      dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java
  13. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java

@ -30,6 +30,8 @@ public interface CalculatedFieldCache {
CalculatedField getCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId);
List<CalculatedField> getCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId);
List<CalculatedFieldLink> getCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId);
List<CalculatedFieldLink> getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId);

37
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java

@ -56,6 +56,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
private final DeviceService deviceService;
private final ConcurrentMap<CalculatedFieldId, CalculatedField> calculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap<EntityId, List<CalculatedField>> entityIdCalculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, List<CalculatedFieldLink>> calculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<EntityId, List<CalculatedFieldLink>> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, CalculatedFieldCtx> calculatedFieldsCtx = new ConcurrentHashMap<>();
@ -65,10 +66,8 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
@Getter
private int initFetchPackSize;
@PostConstruct
public void init() {
// to discuss: fetch on start or fetch on demand
PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize);
cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf));
PageDataIterable<CalculatedFieldLink> cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize);
@ -97,19 +96,37 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
return calculatedField;
}
@Override
public List<CalculatedField> getCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) {
List<CalculatedField> cfs = entityIdCalculatedFields.get(entityId);
if (cfs == null) {
calculatedFieldFetchLock.lock();
try {
cfs = entityIdCalculatedFields.get(entityId);
if (cfs == null) {
cfs = calculatedFieldService.findCalculatedFieldsByEntityId(tenantId, entityId);
entityIdCalculatedFields.put(entityId, cfs);
log.debug("[{}] Fetch calculated fields by entity into cache: {}", entityId, cfs);
}
} finally {
calculatedFieldFetchLock.unlock();
}
}
log.trace("[{}] Found calculated fields by entity in cache: {}", entityId, cfs);
return cfs;
}
@Override
public List<CalculatedFieldLink> getCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId) {
List<CalculatedFieldLink> cfLinks = calculatedFieldLinks.get(calculatedFieldId);
if (cfLinks == null || cfLinks.isEmpty()) {
if (cfLinks == null) {
calculatedFieldFetchLock.lock();
try {
cfLinks = calculatedFieldLinks.get(calculatedFieldId);
if (cfLinks == null || cfLinks.isEmpty()) {
if (cfLinks == null) {
cfLinks = calculatedFieldService.findAllCalculatedFieldLinksById(tenantId, calculatedFieldId);
if (cfLinks != null) {
calculatedFieldLinks.put(calculatedFieldId, cfLinks);
log.debug("[{}] Fetch calculated field links into cache: {}", calculatedFieldId, cfLinks);
}
calculatedFieldLinks.put(calculatedFieldId, cfLinks);
log.debug("[{}] Fetch calculated field links into cache: {}", calculatedFieldId, cfLinks);
}
} finally {
calculatedFieldFetchLock.unlock();
@ -139,7 +156,6 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
return cfLinks;
}
@Override
public void updateCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId) {
log.debug("Update calculated field links per entity for calculated field: [{}]", calculatedFieldId);
@ -225,6 +241,9 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
CalculatedField oldCalculatedField = calculatedFields.remove(calculatedFieldId);
log.debug("[{}] evict calculated field from cache: {}", calculatedFieldId, oldCalculatedField);
calculatedFieldLinks.remove(calculatedFieldId);
log.debug("[{}] evict calculated field from cached calculated fields by entity id: {}", calculatedFieldId, oldCalculatedField);
entityIdCalculatedFields.forEach((entityId, calculatedFields) -> calculatedFields.removeIf(cf -> cf.getId().equals(calculatedFieldId)));
entityIdCalculatedFields.remove(oldCalculatedField.getEntityId());
log.debug("[{}] evict calculated field links from cache: {}", calculatedFieldId, oldCalculatedField);
calculatedFieldsCtx.remove(calculatedFieldId);
log.debug("[{}] evict calculated field ctx from cache: {}", calculatedFieldId, oldCalculatedField);

28
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java

@ -38,6 +38,7 @@ import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
@ -139,6 +140,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field"));
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback"));
scheduledExecutor.submit(() -> rocksDBService.getAll()
.forEach((ctxId, ctx) -> states.put(JacksonUtil.fromString(ctxId, CalculatedFieldEntityCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldEntityCtx.class))));
}
@PreDestroy
@ -337,11 +340,32 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
if (supportedReferencedEntities.contains(entityId.getEntityType())) {
EntityId profileId = getProfileId(tenantId, entityId);
// process by profile
if (profileId != null) {
calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, profileId).forEach(cf -> {
CalculatedFieldLinkConfiguration linkConfiguration = cf.getConfiguration().getReferencedEntityConfig(profileId);
Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(linkConfiguration);
Map<String, KvEntry> updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream()
.filter(entry -> telemetryKeys.containsKey(entry.getKey()))
.collect(Collectors.toMap(
entry -> getMappedKey(entry, telemetryKeys),
entry -> entry,
(v1, v2) -> v1
));
if (!updatedTelemetry.isEmpty()) {
List<CalculatedFieldId> previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds();
executeTelemetryUpdate(tenantId, entityId, cf.getId(), previousCalculatedFieldIds, updatedTelemetry);
}
});
}
// process by links
getCalculatedFieldLinks(tenantId, entityId, profileId).forEach(link -> {
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId();
Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link);
Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link.getConfiguration());
Map<String, KvEntry> updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream()
.filter(entry -> telemetryKeys.containsValue(entry.getKey()))
.filter(entry -> telemetryKeys.containsKey(entry.getKey()))
.collect(Collectors.toMap(
entry -> getMappedKey(entry, telemetryKeys),
entry -> entry,

21
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java

@ -15,10 +15,10 @@
*/
package org.thingsboard.server.service.cf.telemetry;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -28,7 +28,6 @@ import java.util.List;
import java.util.Map;
@Data
@AllArgsConstructor
public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTelemetryUpdateRequest {
private TenantId tenantId;
@ -37,12 +36,20 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel
private List<AttributeKvEntry> kvEntries;
private List<CalculatedFieldId> previousCalculatedFieldIds;
public CalculatedFieldAttributeUpdateRequest(AttributesSaveRequest request) {
this.tenantId = request.getTenantId();
this.entityId = request.getEntityId();
this.scope = request.getScope();
this.kvEntries = request.getEntries();
this.previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds();
}
@Override
public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link) {
public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration) {
return switch (scope) {
case CLIENT_SCOPE -> link.getConfiguration().getClientAttributes();
case SERVER_SCOPE -> link.getConfiguration().getServerAttributes();
case SHARED_SCOPE -> link.getConfiguration().getSharedAttributes();
case CLIENT_SCOPE -> linkConfiguration.getClientAttributes();
case SERVER_SCOPE -> linkConfiguration.getServerAttributes();
case SHARED_SCOPE -> linkConfiguration.getSharedAttributes();
};
}

4
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java

@ -15,7 +15,7 @@
*/
package org.thingsboard.server.service.cf.telemetry;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -34,6 +34,6 @@ public interface CalculatedFieldTelemetryUpdateRequest {
List<CalculatedFieldId> getPreviousCalculatedFieldIds();
Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link);
Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration);
}

16
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java

@ -15,9 +15,9 @@
*/
package org.thingsboard.server.service.cf.telemetry;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -27,7 +27,6 @@ import java.util.List;
import java.util.Map;
@Data
@AllArgsConstructor
public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTelemetryUpdateRequest {
private TenantId tenantId;
@ -35,9 +34,16 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe
private List<TsKvEntry> kvEntries;
private List<CalculatedFieldId> previousCalculatedFieldIds;
public CalculatedFieldTimeSeriesUpdateRequest(TimeseriesSaveRequest request) {
this.tenantId = request.getTenantId();
this.entityId = request.getEntityId();
this.kvEntries = request.getEntries();
this.previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds();
}
@Override
public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link) {
return link.getConfiguration().getTimeSeries();
public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration) {
return linkConfiguration.getTimeSeries();
}
}

7
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -154,9 +154,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
if (request.isSaveLatest() && !request.isOnlyLatest()) {
addEntityViewCallback(tenantId, entityId, request.getEntries());
}
// Use something very similar to addMainCallback. don't forget about tsCallBackExecutor.
//CalculatedFieldTimeSeriesUpdateRequest - add constructor that accepts the TimeseriesSaveRequest
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(tenantId, entityId, request.getEntries(), request.getPreviousCalculatedFieldIds())), tsCallBackExecutor);
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(request)), tsCallBackExecutor);
return saveFuture;
}
@ -172,8 +170,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
addMainCallback(saveFuture, request.getCallback());
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
//CalculatedFieldAttributeUpdateRequest - add constructor that accepts the AttributesSaveRequest
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries(), request.getPreviousCalculatedFieldIds())), tsCallBackExecutor);
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor);
}
@Override

2
common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java

@ -38,6 +38,8 @@ public interface CalculatedFieldService extends EntityDaoService {
List<CalculatedFieldId> findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId);
List<CalculatedField> findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId);
List<CalculatedField> findAllCalculatedFields();
PageData<CalculatedField> findAllCalculatedFields(PageLink pageLink);

17
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java

@ -68,19 +68,22 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel
arguments.entrySet().stream()
.filter(entry -> entry.getValue().getEntityId().equals(entityId))
.forEach(entry -> {
Argument argument = entry.getValue();
Argument tergetArgument = entry.getValue();
String argumentKey = entry.getKey();
switch (argument.getType()) {
switch (tergetArgument.getType()) {
case ATTRIBUTE -> {
switch (argument.getScope()) {
case CLIENT_SCOPE -> linkConfiguration.getClientAttributes().put(entry.getKey(), argument.getKey());
case SERVER_SCOPE -> linkConfiguration.getServerAttributes().put(entry.getKey(), argument.getKey());
case SHARED_SCOPE -> linkConfiguration.getSharedAttributes().put(entry.getKey(), argument.getKey());
switch (tergetArgument.getScope()) {
case CLIENT_SCOPE ->
linkConfiguration.getClientAttributes().put(tergetArgument.getKey(), argumentKey);
case SERVER_SCOPE ->
linkConfiguration.getServerAttributes().put(tergetArgument.getKey(), argumentKey);
case SHARED_SCOPE ->
linkConfiguration.getSharedAttributes().put(tergetArgument.getKey(), argumentKey);
}
}
case TS_LATEST, TS_ROLLING ->
linkConfiguration.getTimeSeries().put(argumentKey, argument.getKey());
linkConfiguration.getTimeSeries().put(tergetArgument.getKey(), argumentKey);
}
});

7
dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

@ -98,6 +98,13 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
return calculatedFieldDao.findCalculatedFieldIdsByEntityId(tenantId, entityId);
}
@Override
public List<CalculatedField> findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) {
log.trace("Executing findCalculatedFieldsByEntityId [{}]", entityId);
validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id);
return calculatedFieldDao.findCalculatedFieldsByEntityId(tenantId, entityId);
}
@Override
public List<CalculatedField> findAllCalculatedFields() {
log.trace("Executing findAll");

2
dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java

@ -31,6 +31,8 @@ public interface CalculatedFieldDao extends Dao<CalculatedField> {
List<CalculatedFieldId> findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId);
List<CalculatedField> findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId);
List<CalculatedField> findAll();
PageData<CalculatedField> findAll(PageLink pageLink);

3
dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.sql.cf;
import org.springframework.data.jpa.repository.JpaRepository;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.dao.model.sql.CalculatedFieldEntity;
@ -28,6 +29,8 @@ public interface CalculatedFieldRepository extends JpaRepository<CalculatedField
List<CalculatedFieldId> findCalculatedFieldIdsByTenantIdAndEntityId(UUID tenantId, UUID entityId);
List<CalculatedField> findAllByTenantIdAndEntityId(UUID tenantId, UUID entityId);
List<CalculatedFieldEntity> findAllByTenantId(UUID tenantId);
List<CalculatedFieldEntity> removeAllByTenantIdAndEntityId(UUID tenantId, UUID entityId);

5
dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java

@ -55,6 +55,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao<CalculatedFieldEntity,
return calculatedFieldRepository.findCalculatedFieldIdsByTenantIdAndEntityId(tenantId.getId(), entityId.getId());
}
@Override
public List<CalculatedField> findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) {
return calculatedFieldRepository.findAllByTenantIdAndEntityId(tenantId.getId(), entityId.getId());
}
@Override
public List<CalculatedField> findAll() {
return DaoUtil.convertDataList(calculatedFieldRepository.findAll());

Loading…
Cancel
Save