Browse Source

added version for relations

pull/10977/head
YevhenBondarenko 2 years ago
parent
commit
358925cffe
  1. 7
      application/src/main/data/upgrade/3.7.0/schema_update.sql
  2. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java
  3. 45
      common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java
  4. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/RelationEntity.java
  5. 112
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  6. 24
      dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java
  7. 157
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java
  8. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationInsertRepository.java
  9. 3
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationRepository.java
  10. 56
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/SqlRelationInsertRepository.java
  11. 3
      dao/src/main/resources/sql/schema-entities.sql
  12. 4
      dao/src/test/java/org/thingsboard/server/dao/service/AlarmServiceTest.java
  13. 92
      dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java
  14. 1
      dao/src/test/resources/sql/psql/drop-all-tables.sql

7
application/src/main/data/upgrade/3.7.0/schema_update.sql

@ -23,3 +23,10 @@ ALTER TABLE attribute_kv ADD COLUMN version bigint default 0;
ALTER TABLE ts_kv_latest ADD COLUMN version bigint default 0;
-- KV VERSIONING UPDATE END
-- RELATION VERSIONING UPDATE START
CREATE SEQUENCE IF NOT EXISTS relation_version_seq cache 1000;
ALTER TABLE relation ADD COLUMN version bigint default 0;
-- RELATION VERSIONING UPDATE END

2
common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java

@ -37,7 +37,7 @@ public interface RelationService {
EntityRelation getRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
boolean saveRelation(TenantId tenantId, EntityRelation relation);
EntityRelation saveRelation(TenantId tenantId, EntityRelation relation);
void saveRelations(TenantId tenantId, List<EntityRelation> relations);

45
common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java

@ -18,8 +18,10 @@ package org.thingsboard.server.common.data.relation;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.databind.JsonNode;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.BaseDataWithAdditionalInfo;
import org.thingsboard.server.common.data.HasVersion;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.validation.Length;
@ -27,7 +29,8 @@ import java.io.Serializable;
@Slf4j
@Schema
public class EntityRelation implements Serializable {
@Data
public class EntityRelation implements HasVersion, Serializable {
private static final long serialVersionUID = 2807343040519543363L;
@ -40,6 +43,7 @@ public class EntityRelation implements Serializable {
@Length(fieldName = "type")
private String type;
private RelationTypeGroup typeGroup;
private Long version;
private transient JsonNode additionalInfo;
@JsonIgnore
private byte[] additionalInfoBytes;
@ -70,6 +74,7 @@ public class EntityRelation implements Serializable {
this.type = entityRelation.getType();
this.typeGroup = entityRelation.getTypeGroup();
this.additionalInfo = entityRelation.getAdditionalInfo();
this.version = entityRelation.getVersion();
}
@Schema(description = "JSON object with [from] Entity Id.", accessMode = Schema.AccessMode.READ_ONLY)
@ -77,35 +82,24 @@ public class EntityRelation implements Serializable {
return from;
}
public void setFrom(EntityId from) {
this.from = from;
}
@Schema(description = "JSON object with [to] Entity Id.", accessMode = Schema.AccessMode.READ_ONLY)
public EntityId getTo() {
return to;
}
public void setTo(EntityId to) {
this.to = to;
}
@Schema(description = "String value of relation type.", example = "Contains")
public String getType() {
return type;
}
public void setType(String type) {
this.type = type;
}
@Schema(description = "Represents the type group of the relation.", example = "COMMON")
public RelationTypeGroup getTypeGroup() {
return typeGroup;
}
public void setTypeGroup(RelationTypeGroup typeGroup) {
this.typeGroup = typeGroup;
@Override
public Long getVersion() {
return version;
}
@Schema(description = "Additional parameters of the relation",implementation = com.fasterxml.jackson.databind.JsonNode.class)
@ -117,25 +111,4 @@ public class EntityRelation implements Serializable {
BaseDataWithAdditionalInfo.setJson(addInfo, json -> this.additionalInfo = json, bytes -> this.additionalInfoBytes = bytes);
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
EntityRelation that = (EntityRelation) o;
if (from != null ? !from.equals(that.from) : that.from != null) return false;
if (to != null ? !to.equals(that.to) : that.to != null) return false;
if (type != null ? !type.equals(that.type) : that.type != null) return false;
return typeGroup == that.typeGroup;
}
@Override
public int hashCode() {
int result = from != null ? from.hashCode() : 0;
result = 31 * result + (to != null ? to.hashCode() : 0);
result = 31 * result + (type != null ? type.hashCode() : 0);
result = 31 * result + (typeGroup != null ? typeGroup.hashCode() : 0);
return result;
}
}

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/RelationEntity.java

@ -23,6 +23,7 @@ import jakarta.persistence.Id;
import jakarta.persistence.IdClass;
import jakarta.persistence.Table;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
@ -40,11 +41,12 @@ import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TO_TYPE_P
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TYPE_GROUP_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TYPE_PROPERTY;
@EqualsAndHashCode(callSuper = true)
@Data
@Entity
@Table(name = RELATION_TABLE_NAME)
@IdClass(RelationCompositeKey.class)
public final class RelationEntity implements ToData<EntityRelation> {
public final class RelationEntity extends VersionedEntity implements ToData<EntityRelation> {
@Id
@Column(name = RELATION_FROM_ID_PROPERTY, columnDefinition = "uuid")
@ -103,6 +105,7 @@ public final class RelationEntity implements ToData<EntityRelation> {
}
relation.setType(relationType);
relation.setTypeGroup(RelationTypeGroup.valueOf(relationTypeGroup));
relation.setVersion(version);
relation.setAdditionalInfo(additionalInfo);
return relation;
}

112
dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java

@ -21,12 +21,13 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.annotation.Lazy;
import org.springframework.dao.ConcurrencyFailureException;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.event.TransactionalEventListener;
@ -52,8 +53,6 @@ import org.thingsboard.server.dao.service.ConstraintValidator;
import org.thingsboard.server.dao.sql.JpaExecutorService;
import org.thingsboard.server.dao.sql.relation.JpaRelationQueryExecutorService;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
@ -149,11 +148,11 @@ public class BaseRelationService implements RelationService {
return relationDao.getRelation(tenantId, from, to, relationType, typeGroup);
},
RelationCacheValue::getRelation,
relations -> RelationCacheValue.builder().relation(relations).build(), false);
relation -> RelationCacheValue.builder().relation(relation).build(), false);
}
@Override
public boolean saveRelation(TenantId tenantId, EntityRelation relation) {
public EntityRelation saveRelation(TenantId tenantId, EntityRelation relation) {
log.trace("Executing saveRelation [{}]", relation);
validate(relation);
var result = relationDao.saveRelation(tenantId, relation);
@ -182,11 +181,11 @@ public class BaseRelationService implements RelationService {
log.trace("Executing saveRelationAsync [{}]", relation);
validate(relation);
var future = relationDao.saveRelationAsync(tenantId, relation);
future.addListener(() -> {
return Futures.transform(future, savedRelation -> {
handleEvictEvent(EntityRelationEvent.from(relation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, relation, ActionType.RELATION_ADD_OR_UPDATE));
return savedRelation != null;
}, MoreExecutors.directExecutor());
return future;
}
@Override
@ -194,10 +193,11 @@ public class BaseRelationService implements RelationService {
log.trace("Executing DeleteRelation [{}]", relation);
validate(relation);
var result = relationDao.deleteRelation(tenantId, relation);
//TODO: evict cache only if the relation was deleted. Note: relationDao.deleteRelation requires improvement.
publishEvictEvent(EntityRelationEvent.from(relation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, relation, ActionType.RELATION_DELETED));
return result;
if (result != null) {
publishEvictEvent(EntityRelationEvent.from(relation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, relation, ActionType.RELATION_DELETED));
}
return result != null;
}
@Override
@ -205,11 +205,13 @@ public class BaseRelationService implements RelationService {
log.trace("Executing deleteRelationAsync [{}]", relation);
validate(relation);
var future = relationDao.deleteRelationAsync(tenantId, relation);
future.addListener(() -> {
handleEvictEvent(EntityRelationEvent.from(relation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, relation, ActionType.RELATION_DELETED));
return Futures.transform(future, deletedRelation -> {
if (deletedRelation != null) {
handleEvictEvent(EntityRelationEvent.from(relation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, relation, ActionType.RELATION_DELETED));
}
return deletedRelation != null;
}, MoreExecutors.directExecutor());
return future;
}
@Override
@ -217,11 +219,11 @@ public class BaseRelationService implements RelationService {
log.trace("Executing deleteRelation [{}][{}][{}][{}]", from, to, relationType, typeGroup);
validate(from, to, relationType, typeGroup);
var result = relationDao.deleteRelation(tenantId, from, to, relationType, typeGroup);
//TODO: evict cache only if the relation was deleted. Note: relationDao.deleteRelation requires improvement.
EntityRelation entityRelation = new EntityRelation(from, to, relationType, typeGroup);
publishEvictEvent(EntityRelationEvent.from(entityRelation));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, entityRelation, ActionType.RELATION_DELETED));
return result;
if (result != null) {
publishEvictEvent(EntityRelationEvent.from(result));
eventPublisher.publishEvent(new RelationActionEvent(tenantId, result, ActionType.RELATION_DELETED));
}
return result != null;
}
@Override
@ -229,9 +231,12 @@ public class BaseRelationService implements RelationService {
log.trace("Executing deleteRelationAsync [{}][{}][{}][{}]", from, to, relationType, typeGroup);
validate(from, to, relationType, typeGroup);
var future = relationDao.deleteRelationAsync(tenantId, from, to, relationType, typeGroup);
EntityRelationEvent event = new EntityRelationEvent(from, to, relationType, typeGroup);
future.addListener(() -> handleEvictEvent(event), MoreExecutors.directExecutor());
return future;
return Futures.transform(future, deletedEvent -> {
if (deletedEvent != null) {
handleEvictEvent(EntityRelationEvent.from(deletedEvent));
}
return deletedEvent != null;
}, MoreExecutors.directExecutor());
}
@Transactional
@ -250,60 +255,27 @@ public class BaseRelationService implements RelationService {
public void deleteEntityRelations(TenantId tenantId, EntityId entityId, RelationTypeGroup relationTypeGroup) {
log.trace("Executing deleteEntityRelations [{}]", entityId);
validate(entityId);
List<EntityRelation> inboundRelations = relationTypeGroup == null
? relationDao.findAllByTo(tenantId, entityId)
: relationDao.findAllByTo(tenantId, entityId, relationTypeGroup);
List<EntityRelation> outboundRelations = relationTypeGroup == null
? relationDao.findAllByFrom(tenantId, entityId)
: relationDao.findAllByFrom(tenantId, entityId, relationTypeGroup);
if (!inboundRelations.isEmpty()) {
try {
if (relationTypeGroup == null) {
relationDao.deleteInboundRelations(tenantId, entityId);
} else {
relationDao.deleteInboundRelations(tenantId, entityId, relationTypeGroup);
}
} catch (ConcurrencyFailureException e) {
log.debug("Concurrency exception while deleting relations [{}]", inboundRelations, e);
}
for (EntityRelation relation : inboundRelations) {
eventPublisher.publishEvent(EntityRelationEvent.from(relation));
}
List<EntityRelation> inboundRelations;
if (relationTypeGroup == null) {
inboundRelations = relationDao.deleteInboundRelations(tenantId, entityId);
} else {
inboundRelations = relationDao.deleteInboundRelations(tenantId, entityId, relationTypeGroup);
}
if (!outboundRelations.isEmpty()) {
if (relationTypeGroup == null) {
relationDao.deleteOutboundRelations(tenantId, entityId);
} else {
relationDao.deleteOutboundRelations(tenantId, entityId, relationTypeGroup);
}
for (EntityRelation relation : outboundRelations) {
eventPublisher.publishEvent(EntityRelationEvent.from(relation));
}
for (EntityRelation relation : inboundRelations) {
eventPublisher.publishEvent(EntityRelationEvent.from(relation));
}
}
private List<ListenableFuture<Boolean>> deleteRelationGroupsAsync(TenantId tenantId, List<List<EntityRelation>> relations, boolean deleteFromDb) {
List<ListenableFuture<Boolean>> results = new ArrayList<>();
for (List<EntityRelation> relationList : relations) {
relationList.forEach(relation -> results.add(deleteAsync(tenantId, relation, deleteFromDb)));
List<EntityRelation> outboundRelations;
if (relationTypeGroup == null) {
outboundRelations = relationDao.deleteOutboundRelations(tenantId, entityId);
} else {
outboundRelations = relationDao.deleteOutboundRelations(tenantId, entityId, relationTypeGroup);
}
return results;
}
private ListenableFuture<Boolean> deleteAsync(TenantId tenantId, EntityRelation relation, boolean deleteFromDb) {
if (deleteFromDb) {
return Futures.transform(relationDao.deleteRelationAsync(tenantId, relation),
bool -> {
handleEvictEvent(EntityRelationEvent.from(relation));
return bool;
}, MoreExecutors.directExecutor());
} else {
handleEvictEvent(EntityRelationEvent.from(relation));
return Futures.immediateFuture(false);
for (EntityRelation relation : outboundRelations) {
eventPublisher.publishEvent(EntityRelationEvent.from(relation));
}
}

24
dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java

@ -48,29 +48,27 @@ public interface RelationDao {
EntityRelation getRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
boolean saveRelation(TenantId tenantId, EntityRelation relation);
EntityRelation saveRelation(TenantId tenantId, EntityRelation relation);
void saveRelations(TenantId tenantId, Collection<EntityRelation> relations);
List<EntityRelation> saveRelations(TenantId tenantId, List<EntityRelation> relations);
ListenableFuture<Boolean> saveRelationAsync(TenantId tenantId, EntityRelation relation);
ListenableFuture<EntityRelation> saveRelationAsync(TenantId tenantId, EntityRelation relation);
boolean deleteRelation(TenantId tenantId, EntityRelation relation);
EntityRelation deleteRelation(TenantId tenantId, EntityRelation relation);
ListenableFuture<Boolean> deleteRelationAsync(TenantId tenantId, EntityRelation relation);
ListenableFuture<EntityRelation> deleteRelationAsync(TenantId tenantId, EntityRelation relation);
boolean deleteRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
EntityRelation deleteRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
ListenableFuture<Boolean> deleteRelationAsync(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
ListenableFuture<EntityRelation> deleteRelationAsync(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup);
void deleteOutboundRelations(TenantId tenantId, EntityId entity);
List<EntityRelation> deleteOutboundRelations(TenantId tenantId, EntityId entity);
void deleteOutboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup);
List<EntityRelation> deleteOutboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup);
void deleteInboundRelations(TenantId tenantId, EntityId entity);
List<EntityRelation> deleteInboundRelations(TenantId tenantId, EntityId entity);
void deleteInboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup);
ListenableFuture<Boolean> deleteOutboundRelationsAsync(TenantId tenantId, EntityId entity);
List<EntityRelation> deleteInboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup);
List<EntityRelation> findRuleNodeToRuleChainRelations(RuleChainType ruleChainType, int limit);
}

157
dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java

@ -18,11 +18,11 @@ package org.thingsboard.server.dao.sql.relation;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.dao.ConcurrencyFailureException;
import org.springframework.dao.DataAccessException;
import org.springframework.data.domain.PageRequest;
import org.springframework.stereotype.Component;
import org.springframework.util.CollectionUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
@ -36,11 +36,19 @@ import org.thingsboard.server.dao.util.SqlDao;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_FROM_ID_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_FROM_TYPE_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TO_ID_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TO_TYPE_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TYPE_GROUP_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TYPE_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.VERSION_COLUMN;
/**
* Created by Valerii Sosliuk on 5/29/2017.
*/
@ -50,6 +58,8 @@ import java.util.stream.Collectors;
public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService implements RelationDao {
private static final List<String> ALL_TYPE_GROUP_NAMES = new ArrayList<>();
private static final String RETURNING = "RETURNING from_id, from_type, to_id, to_type, relation_type, relation_type_group, nextval('relation_version_seq') as version";
private static final String DELETE_QUERY = "DELETE FROM relation WHERE from_id = ? AND from_type = ? AND to_id = ? AND to_type = ? AND relation_type = ? AND relation_type_group = ? " + RETURNING;
static {
Arrays.stream(RelationTypeGroup.values()).map(RelationTypeGroup::name).forEach(ALL_TYPE_GROUP_NAMES::add);
@ -144,107 +154,138 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple
}
@Override
public boolean saveRelation(TenantId tenantId, EntityRelation relation) {
return relationInsertRepository.saveOrUpdate(new RelationEntity(relation)) != null;
public EntityRelation saveRelation(TenantId tenantId, EntityRelation relation) {
return DaoUtil.getData(relationInsertRepository.saveOrUpdate(new RelationEntity(relation)));
}
@Override
public void saveRelations(TenantId tenantId, Collection<EntityRelation> relations) {
public List<EntityRelation> saveRelations(TenantId tenantId, List<EntityRelation> relations) {
List<RelationEntity> entities = relations.stream().map(RelationEntity::new).collect(Collectors.toList());
relationInsertRepository.saveOrUpdate(entities);
return DaoUtil.convertDataList(relationInsertRepository.saveOrUpdate(entities));
}
@Override
public ListenableFuture<Boolean> saveRelationAsync(TenantId tenantId, EntityRelation relation) {
return service.submit(() -> relationInsertRepository.saveOrUpdate(new RelationEntity(relation)) != null);
public ListenableFuture<EntityRelation> saveRelationAsync(TenantId tenantId, EntityRelation relation) {
return service.submit(() -> DaoUtil.getData(relationInsertRepository.saveOrUpdate(new RelationEntity(relation))));
}
@Override
public boolean deleteRelation(TenantId tenantId, EntityRelation relation) {
public EntityRelation deleteRelation(TenantId tenantId, EntityRelation relation) {
RelationCompositeKey key = new RelationCompositeKey(relation);
return deleteRelationIfExists(key);
}
@Override
public ListenableFuture<Boolean> deleteRelationAsync(TenantId tenantId, EntityRelation relation) {
public ListenableFuture<EntityRelation> deleteRelationAsync(TenantId tenantId, EntityRelation relation) {
RelationCompositeKey key = new RelationCompositeKey(relation);
return service.submit(
() -> deleteRelationIfExists(key));
}
@Override
public boolean deleteRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) {
public EntityRelation deleteRelation(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) {
RelationCompositeKey key = getRelationCompositeKey(from, to, relationType, typeGroup);
return deleteRelationIfExists(key);
}
@Override
public ListenableFuture<Boolean> deleteRelationAsync(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) {
public ListenableFuture<EntityRelation> deleteRelationAsync(TenantId tenantId, EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) {
RelationCompositeKey key = getRelationCompositeKey(from, to, relationType, typeGroup);
return service.submit(
() -> deleteRelationIfExists(key));
}
private boolean deleteRelationIfExists(RelationCompositeKey key) {
boolean relationExistsBeforeDelete = relationRepository.existsById(key);
if (relationExistsBeforeDelete) {
try {
relationRepository.deleteById(key);
} catch (DataAccessException e) {
log.debug("[{}] Concurrency exception while deleting relation", key, e);
private EntityRelation deleteRelationIfExists(RelationCompositeKey key) {
return jdbcTemplate.query(DELETE_QUERY, rs -> {
if (!rs.next()) {
return null;
}
}
return relationExistsBeforeDelete;
EntityRelation relation = new EntityRelation();
var fromId = rs.getObject(RELATION_FROM_ID_PROPERTY, UUID.class);
var fromType = rs.getString(RELATION_FROM_TYPE_PROPERTY);
var toId = rs.getObject(RELATION_TO_ID_PROPERTY, UUID.class);
var toType = rs.getString(RELATION_TO_TYPE_PROPERTY);
var relationTypeGroup = rs.getString(RELATION_TYPE_GROUP_PROPERTY);
var relationType = rs.getString(RELATION_TYPE_PROPERTY);
var version = rs.getLong(VERSION_COLUMN);
//additionalInfo ignored (no need to send extra data for delete events)
relation.setTo(EntityIdFactory.getByTypeAndUuid(toType, toId));
relation.setFrom(EntityIdFactory.getByTypeAndUuid(fromType, fromId));
relation.setType(relationType);
relation.setTypeGroup(RelationTypeGroup.valueOf(relationTypeGroup));
relation.setVersion(version);
return relation;
}, key.getFromId(), key.getFromType(), key.getToId(), key.getToType(), key.getRelationType(), key.getRelationTypeGroup());
}
@Override
public void deleteOutboundRelations(TenantId tenantId, EntityId entity) {
try {
relationRepository.deleteByFromIdAndFromType(entity.getId(), entity.getEntityType().name());
} catch (ConcurrencyFailureException e) {
log.debug("Concurrency exception while deleting relations [{}]", entity, e);
}
public List<EntityRelation> deleteOutboundRelations(TenantId tenantId, EntityId entity) {
return deleteRelations(entity, null, false);
}
@Override
public void deleteOutboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup) {
try {
relationRepository.deleteByFromIdAndFromTypeAndRelationTypeGroupIn(entity.getId(), entity.getEntityType().name(), Collections.singletonList(relationTypeGroup.name()));
} catch (ConcurrencyFailureException e) {
log.debug("Concurrency exception while deleting relations [{}]", entity, e);
}
public List<EntityRelation> deleteOutboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup) {
return deleteRelations(entity, Collections.singletonList(relationTypeGroup.name()), false);
}
@Override
public void deleteInboundRelations(TenantId tenantId, EntityId entity) {
try {
relationRepository.deleteByToIdAndToTypeAndRelationTypeGroupIn(entity.getId(), entity.getEntityType().name(), ALL_TYPE_GROUP_NAMES);
} catch (ConcurrencyFailureException e) {
log.debug("Concurrency exception while deleting relations [{}]", entity, e);
}
public List<EntityRelation> deleteInboundRelations(TenantId tenantId, EntityId entity) {
return deleteRelations(entity, ALL_TYPE_GROUP_NAMES, true);
}
@Override
public void deleteInboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup) {
try {
relationRepository.deleteByToIdAndToTypeAndRelationTypeGroupIn(entity.getId(), entity.getEntityType().name(), Collections.singletonList(relationTypeGroup.name()));
} catch (ConcurrencyFailureException e) {
log.debug("Concurrency exception while deleting relations [{}]", entity, e);
}
public List<EntityRelation> deleteInboundRelations(TenantId tenantId, EntityId entity, RelationTypeGroup relationTypeGroup) {
return deleteRelations(entity, Collections.singletonList(relationTypeGroup.name()), true);
}
@Override
public ListenableFuture<Boolean> deleteOutboundRelationsAsync(TenantId tenantId, EntityId entity) {
return service.submit(
() -> {
boolean relationExistsBeforeDelete = relationRepository
.findAllByFromIdAndFromType(entity.getId(), entity.getEntityType().name())
.size() > 0;
if (relationExistsBeforeDelete) {
relationRepository.deleteByFromIdAndFromType(entity.getId(), entity.getEntityType().name());
}
return relationExistsBeforeDelete;
});
private List<EntityRelation> deleteRelations(EntityId entityId, List<String> relationTypeGroups, boolean inbound) {
List<Object> params = new ArrayList<>();
params.add(entityId.getId());
params.add(entityId.getEntityType().name());
StringBuilder sqlBuilder = new StringBuilder("DELETE FROM relation WHERE ");
if (inbound) {
sqlBuilder.append("to_id = ? AND to_type = ? ");
} else {
sqlBuilder.append("from_id = ? AND from_type = ? ");
}
if (!CollectionUtils.isEmpty(relationTypeGroups)) {
sqlBuilder.append("AND relation_type_group IN (?");
for (int i = 1; i < relationTypeGroups.size(); i++) {
sqlBuilder.append(", ?");
}
sqlBuilder.append(")");
params.addAll(relationTypeGroups);
}
sqlBuilder.append(RETURNING);
return jdbcTemplate.queryForList(sqlBuilder.toString(), params.toArray()).stream()
.map(row -> {
EntityRelation relation = new EntityRelation();
var fromId = row.get(RELATION_FROM_ID_PROPERTY);
var fromType = row.get(RELATION_FROM_TYPE_PROPERTY);
var toId = row.get(RELATION_TO_ID_PROPERTY);
var toType = row.get(RELATION_TO_TYPE_PROPERTY);
var relationTypeGroup = row.get(RELATION_TYPE_GROUP_PROPERTY);
var relationType = row.get(RELATION_TYPE_PROPERTY);
var version = row.get(VERSION_COLUMN);
//additionalInfo ignored (no need to send extra data for delete events)
relation.setTo(EntityIdFactory.getByTypeAndUuid((String) toType, (UUID) toId));
relation.setFrom(EntityIdFactory.getByTypeAndUuid((String) fromType, (UUID) fromId));
relation.setType((String) relationType);
relation.setTypeGroup(RelationTypeGroup.valueOf((String) relationTypeGroup));
relation.setVersion((Long) version);
return relation;
})
.collect(Collectors.toList());
}
@Override

2
dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationInsertRepository.java

@ -23,6 +23,6 @@ public interface RelationInsertRepository {
RelationEntity saveOrUpdate(RelationEntity entity);
void saveOrUpdate(List<RelationEntity> entities);
List<RelationEntity> saveOrUpdate(List<RelationEntity> entities);
}

3
dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationRepository.java

@ -58,9 +58,6 @@ public interface RelationRepository
String relationType,
String relationTypeGroup);
List<RelationEntity> findAllByFromIdAndFromType(UUID fromId,
String fromType);
@Query("SELECT r FROM RelationEntity r WHERE " +
"r.relationTypeGroup = 'RULE_NODE' AND r.toType = 'RULE_CHAIN' " +
"AND r.toId in (SELECT id from RuleChainEntity where type = :ruleChainType )")

56
dao/src/main/java/org/thingsboard/server/dao/sql/relation/SqlRelationInsertRepository.java

@ -15,33 +15,39 @@
*/
package org.thingsboard.server.dao.sql.relation;
import jakarta.persistence.EntityManager;
import jakarta.persistence.PersistenceContext;
import jakarta.persistence.Query;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.PreparedStatementCreator;
import org.springframework.jdbc.core.SqlProvider;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.dao.model.sql.RelationEntity;
import jakarta.persistence.EntityManager;
import jakarta.persistence.PersistenceContext;
import jakarta.persistence.Query;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.List;
import static org.thingsboard.server.dao.model.ModelConstants.VERSION_COLUMN;
@Repository
@Transactional
public class SqlRelationInsertRepository implements RelationInsertRepository {
private static final String INSERT_ON_CONFLICT_DO_UPDATE_JPA = "INSERT INTO relation (from_id, from_type, to_id, to_type, relation_type_group, relation_type, additional_info)" +
" VALUES (:fromId, :fromType, :toId, :toType, :relationTypeGroup, :relationType, :additionalInfo) " +
"ON CONFLICT (from_id, from_type, relation_type_group, relation_type, to_id, to_type) DO UPDATE SET additional_info = :additionalInfo returning *";
private static final String INSERT_ON_CONFLICT_DO_UPDATE_JDBC = "INSERT INTO relation (from_id, from_type, to_id, to_type, relation_type_group, relation_type, additional_info)" +
" VALUES (?, ?, ?, ?, ?, ?, ?) " +
"ON CONFLICT (from_id, from_type, relation_type_group, relation_type, to_id, to_type) DO UPDATE SET additional_info = ?";
private static final String INSERT_ON_CONFLICT_DO_UPDATE_JPA = "INSERT INTO relation (from_id, from_type, to_id, to_type, relation_type_group, relation_type, version, additional_info)" +
" VALUES (:fromId, :fromType, :toId, :toType, :relationTypeGroup, :relationType, nextval('relation_version_seq'), :additionalInfo) " +
"ON CONFLICT (from_id, from_type, relation_type_group, relation_type, to_id, to_type) DO UPDATE SET additional_info = :additionalInfo, version = nextval('relation_version_seq') returning *";
private static final String INSERT_ON_CONFLICT_DO_UPDATE_JDBC = "INSERT INTO relation (from_id, from_type, to_id, to_type, relation_type_group, relation_type, version, additional_info)" +
" VALUES (?, ?, ?, ?, ?, ?, nextval('relation_version_seq'), ?) " +
"ON CONFLICT (from_id, from_type, relation_type_group, relation_type, to_id, to_type) DO UPDATE SET additional_info = ?, version = nextval('relation_version_seq')";
@PersistenceContext
protected EntityManager entityManager;
@ -71,8 +77,9 @@ public class SqlRelationInsertRepository implements RelationInsertRepository {
}
@Override
public void saveOrUpdate(List<RelationEntity> entities) {
jdbcTemplate.batchUpdate(INSERT_ON_CONFLICT_DO_UPDATE_JDBC, new BatchPreparedStatementSetter() {
public List<RelationEntity> saveOrUpdate(List<RelationEntity> entities) {
KeyHolder keyHolder = new GeneratedKeyHolder();
jdbcTemplate.batchUpdate(new SequencePreparedStatementCreator(INSERT_ON_CONFLICT_DO_UPDATE_JDBC), new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
RelationEntity relation = entities.get(i);
@ -98,7 +105,30 @@ public class SqlRelationInsertRepository implements RelationInsertRepository {
public int getBatchSize() {
return entities.size();
}
});
}, keyHolder);
var seqNumbers = keyHolder.getKeyList();
for (int i = 0; i < entities.size(); i++) {
entities.get(i).setVersion((Long) seqNumbers.get(i).get(VERSION_COLUMN));
}
return entities;
}
private record SequencePreparedStatementCreator(String sql) implements PreparedStatementCreator, SqlProvider {
private static final String[] COLUMNS = {VERSION_COLUMN};
@Override
public PreparedStatement createPreparedStatement(Connection con) throws SQLException {
return con.prepareStatement(sql, COLUMNS);
}
@Override
public String getSql() {
return this.sql;
}
}
}

3
dao/src/main/resources/sql/schema-entities.sql

@ -420,6 +420,8 @@ CREATE TABLE IF NOT EXISTS error_event (
e_error varchar
) PARTITION BY RANGE (ts);
CREATE SEQUENCE IF NOT EXISTS relation_version_seq cache 1000;
CREATE TABLE IF NOT EXISTS relation (
from_id uuid,
from_type varchar(255),
@ -428,6 +430,7 @@ CREATE TABLE IF NOT EXISTS relation (
relation_type_group varchar(255),
relation_type varchar(255),
additional_info varchar,
version bigint default 0,
CONSTRAINT relation_pkey PRIMARY KEY (from_id, from_type, relation_type_group, relation_type, to_id, to_type)
);

4
dao/src/test/java/org/thingsboard/server/dao/service/AlarmServiceTest.java

@ -340,7 +340,7 @@ public class AlarmServiceTest extends AbstractServiceTest {
EntityRelation relation = new EntityRelation(parentId, childId, EntityRelation.CONTAINS_TYPE);
Assert.assertTrue(relationService.saveRelation(tenantId, relation));
Assert.assertNotNull(relationService.saveRelation(tenantId, relation));
long ts = System.currentTimeMillis();
AlarmApiCallResult result = alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.builder()
@ -877,7 +877,7 @@ public class AlarmServiceTest extends AbstractServiceTest {
EntityRelation relation = new EntityRelation(parentId, childId, EntityRelation.CONTAINS_TYPE);
Assert.assertTrue(relationService.saveRelation(tenantId, relation));
Assert.assertNotNull(relationService.saveRelation(tenantId, relation));
long ts = System.currentTimeMillis();
AlarmApiCallResult result = alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.builder()

92
dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java

@ -57,13 +57,13 @@ public class RelationServiceTest extends AbstractServiceTest {
}
@Test
public void testSaveRelation() throws ExecutionException, InterruptedException {
public void testSaveRelation() {
AssetId parentId = new AssetId(Uuids.timeBased());
AssetId childId = new AssetId(Uuids.timeBased());
EntityRelation relation = new EntityRelation(parentId, childId, EntityRelation.CONTAINS_TYPE);
Assert.assertTrue(saveRelation(relation));
Assert.assertNotNull(saveRelation(relation));
Assert.assertTrue(relationService.checkRelation(SYSTEM_TENANT_ID, parentId, childId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.COMMON));
@ -204,8 +204,8 @@ public class RelationServiceTest extends AbstractServiceTest {
Assert.assertEquals(0, relations.size());
}
private Boolean saveRelation(EntityRelation relationA1) {
return relationService.saveRelation(SYSTEM_TENANT_ID, relationA1);
private EntityRelation saveRelation(EntityRelation relation) {
return relationService.saveRelation(SYSTEM_TENANT_ID, relation);
}
@Test
@ -265,9 +265,9 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationB = new EntityRelation(assetB, assetC, EntityRelation.CONTAINS_TYPE);
EntityRelation relationC = new EntityRelation(assetC, assetA, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, -1, false));
@ -299,8 +299,8 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationBD = new EntityRelation(assetB, deviceD, EntityRelation.CONTAINS_TYPE);
saveRelation(relationAB);
saveRelation(relationBC);
relationAB = saveRelation(relationAB);
relationBC = saveRelation(relationBC);
saveRelation(relationBD);
EntityRelationsQuery query = new EntityRelationsQuery();
@ -329,26 +329,20 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationAB = new EntityRelation(root, left, EntityRelation.CONTAINS_TYPE);
EntityRelation relationBC = new EntityRelation(root, right, EntityRelation.CONTAINS_TYPE);
saveRelation(relationAB);
expected.add(relationAB);
saveRelation(relationBC);
expected.add(relationBC);
expected.add(saveRelation(relationAB));
expected.add(saveRelation(relationBC));
for (int i = 0; i < maxLevel; i++) {
var newLeft = new AssetId(Uuids.timeBased());
var newRight = new AssetId(Uuids.timeBased());
EntityRelation relationLeft = new EntityRelation(left, newLeft, EntityRelation.CONTAINS_TYPE);
EntityRelation relationRight = new EntityRelation(right, newRight, EntityRelation.CONTAINS_TYPE);
saveRelation(relationLeft);
expected.add(relationLeft);
saveRelation(relationRight);
expected.add(relationRight);
expected.add(saveRelation(relationLeft));
expected.add(saveRelation(relationRight));
left = newLeft;
right = newRight;
}
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(root, EntitySearchDirection.FROM, -1, false));
query.setFilters(Collections.singletonList(new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.singletonList(EntityType.ASSET))));
@ -372,7 +366,7 @@ public class RelationServiceTest extends AbstractServiceTest {
relation.setTo(new AssetId(Uuids.timeBased()));
relation.setType(EntityRelation.CONTAINS_TYPE);
Assertions.assertThrows(DataValidationException.class, () -> {
Assert.assertTrue(saveRelation(relation));
Assert.assertNotNull(saveRelation(relation));
});
}
@ -382,7 +376,7 @@ public class RelationServiceTest extends AbstractServiceTest {
relation.setFrom(new AssetId(Uuids.timeBased()));
relation.setType(EntityRelation.CONTAINS_TYPE);
Assertions.assertThrows(DataValidationException.class, () -> {
Assert.assertTrue(saveRelation(relation));
Assert.assertNotNull(saveRelation(relation));
});
}
@ -392,7 +386,7 @@ public class RelationServiceTest extends AbstractServiceTest {
relation.setFrom(new AssetId(Uuids.timeBased()));
relation.setTo(new AssetId(Uuids.timeBased()));
Assertions.assertThrows(DataValidationException.class, () -> {
Assert.assertTrue(saveRelation(relation));
Assert.assertNotNull(saveRelation(relation));
});
}
@ -414,10 +408,10 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationC = new EntityRelation(assetC, assetD, EntityRelation.CONTAINS_TYPE);
EntityRelation relationD = new EntityRelation(assetC, assetE, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
saveRelation(relationD);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
relationD = saveRelation(relationD);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, -1, true));
@ -450,9 +444,9 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationB = new EntityRelation(assetB, assetC, EntityRelation.CONTAINS_TYPE);
EntityRelation relationC = new EntityRelation(assetC, assetD, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, -1, true));
@ -494,12 +488,12 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationE = new EntityRelation(assetD, assetF, EntityRelation.CONTAINS_TYPE);
EntityRelation relationF = new EntityRelation(assetD, assetG, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
saveRelation(relationD);
saveRelation(relationE);
saveRelation(relationF);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
relationD = saveRelation(relationD);
relationE = saveRelation(relationE);
relationF = saveRelation(relationF);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, 2, true));
@ -547,12 +541,12 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationE = new EntityRelation(assetD, assetF, EntityRelation.CONTAINS_TYPE);
EntityRelation relationF = new EntityRelation(assetD, assetG, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
saveRelation(relationD);
saveRelation(relationE);
saveRelation(relationF);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
relationD = saveRelation(relationD);
relationE = saveRelation(relationE);
relationF = saveRelation(relationF);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, 2, false));
@ -600,12 +594,12 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation relationE = new EntityRelation(assetD, assetF, EntityRelation.CONTAINS_TYPE);
EntityRelation relationF = new EntityRelation(assetD, assetG, EntityRelation.CONTAINS_TYPE);
saveRelation(relationA);
saveRelation(relationB);
saveRelation(relationC);
saveRelation(relationD);
saveRelation(relationE);
saveRelation(relationF);
relationA = saveRelation(relationA);
relationB = saveRelation(relationB);
relationC = saveRelation(relationC);
relationD = saveRelation(relationD);
relationE = saveRelation(relationE);
relationF = saveRelation(relationF);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(assetA, EntitySearchDirection.FROM, -1, false));
@ -670,8 +664,8 @@ public class RelationServiceTest extends AbstractServiceTest {
EntityRelation firstRelation = new EntityRelation(rootAsset, firstAsset, EntityRelation.CONTAINS_TYPE);
EntityRelation secondRelation = new EntityRelation(rootAsset, secondAsset, EntityRelation.CONTAINS_TYPE);
saveRelation(firstRelation);
saveRelation(secondRelation);
firstRelation = saveRelation(firstRelation);
secondRelation = saveRelation(secondRelation);
if (!lastLvlOnly || lvl == 1) {
entityRelations.add(firstRelation);

1
dao/src/test/resources/sql/psql/drop-all-tables.sql

@ -34,6 +34,7 @@ DROP TABLE IF EXISTS stats_event;
DROP TABLE IF EXISTS lc_event;
DROP TABLE IF EXISTS error_event;
DROP TABLE IF EXISTS relation;
DROP SEQUENCE IF EXISTS relation_version_seq;
DROP TABLE IF EXISTS tenant;
DROP TABLE IF EXISTS ts_kv;
DROP TABLE IF EXISTS ts_kv_latest;

Loading…
Cancel
Save