Browse Source

Refactor daos for partitioned entities

pull/10083/head
ViacheslavKlimov 3 years ago
parent
commit
90f971d018
  1. 6
      dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmCommentDao.java
  2. 2
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java
  3. 2
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogDao.java
  4. 31
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
  5. 54
      dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java
  6. 27
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmCommentDao.java
  7. 24
      dao/src/main/java/org/thingsboard/server/dao/sql/audit/JpaAuditLogDao.java
  8. 17
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java
  9. 21
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java
  10. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageInfoDao.java
  11. 2
      dao/src/test/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmCommentDaoTest.java

6
dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmCommentDao.java

@ -18,7 +18,6 @@ package org.thingsboard.server.dao.alarm;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.alarm.AlarmComment;
import org.thingsboard.server.common.data.alarm.AlarmCommentInfo;
import org.thingsboard.server.common.data.id.AlarmCommentId;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
@ -29,13 +28,10 @@ import java.util.UUID;
public interface AlarmCommentDao extends Dao<AlarmComment> {
AlarmComment createAlarmComment(TenantId tenantId, AlarmComment alarmComment);
void deleteAlarmComment(TenantId tenantId, AlarmCommentId alarmCommentId);
AlarmComment findAlarmCommentById(TenantId tenantId, UUID key);
PageData<AlarmCommentInfo> findAlarmComments(TenantId tenantId, AlarmId id, PageLink pageLink);
ListenableFuture<AlarmComment> findAlarmCommentByIdAsync(TenantId tenantId, UUID key);
}

2
dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java

@ -86,7 +86,7 @@ public class BaseAlarmCommentService extends AbstractEntityService implements Al
if (alarmComment.getType() == null) {
alarmComment.setType(AlarmCommentType.OTHER);
}
return alarmCommentDao.createAlarmComment(tenantId, alarmComment);
return alarmCommentDao.save(tenantId, alarmComment);
}
private AlarmComment updateAlarmComment(TenantId tenantId, AlarmComment newAlarmComment) {

2
dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogDao.java

@ -30,8 +30,6 @@ import java.util.UUID;
public interface AuditLogDao extends Dao<AuditLog> {
ListenableFuture<AuditLog> saveByTenantId(AuditLog auditLog);
PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink);
PageData<AuditLog> findAuditLogsByTenantIdAndCustomerId(UUID tenantId, CustomerId customerId, List<ActionType> actionTypes, TimePageLink pageLink);

31
dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java

@ -20,7 +20,6 @@ import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
@ -48,6 +47,7 @@ import org.thingsboard.server.dao.audit.sink.AuditLogSink;
import org.thingsboard.server.dao.device.provision.ProvisionRequest;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.sql.JpaExecutorService;
import java.io.PrintWriter;
import java.io.StringWriter;
@ -76,6 +76,9 @@ public class AuditLogServiceImpl implements AuditLogService {
@Autowired
private AuditLogSink auditLogSink;
@Autowired
private JpaExecutorService executor;
@Autowired
private DataValidator<AuditLog> auditLogValidator;
@ -380,15 +383,15 @@ public class AuditLogServiceImpl implements AuditLogService {
}
private ListenableFuture<Void> logAction(TenantId tenantId,
EntityId entityId,
String entityName,
CustomerId customerId,
UserId userId,
String userName,
ActionType actionType,
JsonNode actionData,
ActionStatus actionStatus,
String actionFailureDetails) {
EntityId entityId,
String entityName,
CustomerId customerId,
UserId userId,
String userName,
ActionType actionType,
JsonNode actionData,
ActionStatus actionStatus,
String actionFailureDetails) {
AuditLog auditLogEntry = createAuditLogEntry(tenantId, entityId, entityName, customerId, userId, userName,
actionType, actionData, actionStatus, actionFailureDetails);
log.trace("Executing logAction [{}]", auditLogEntry);
@ -402,11 +405,11 @@ public class AuditLogServiceImpl implements AuditLogService {
}
}
ListenableFuture<AuditLog> future = auditLogDao.saveByTenantId(auditLogEntry);
return Futures.transform(future, auditLog -> {
auditLogSink.logAction(auditLogEntry);
return executor.submit(() -> {
AuditLog auditLog = auditLogDao.save(tenantId, auditLogEntry);
auditLogSink.logAction(auditLog);
return null;
}, MoreExecutors.directExecutor());
});
}
}

54
dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java

@ -19,12 +19,8 @@ import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.Dao;
import org.thingsboard.server.dao.DaoUtil;
@ -37,45 +33,26 @@ import java.util.Collection;
import java.util.List;
import java.util.Optional;
import java.util.UUID;
import java.util.function.Consumer;
/**
* @author Valerii Sosliuk
*/
@Slf4j
@SqlDao
public abstract class JpaAbstractDao<E extends BaseEntity<D>, D extends HasId<?>>
public abstract class JpaAbstractDao<E extends BaseEntity<D>, D>
extends JpaAbstractDaoListeningExecutorService
implements Dao<D> {
@PersistenceContext
private EntityManager entityManager;
@Autowired
private TransactionTemplate transactionTemplate;
protected abstract Class<E> getEntityClass();
protected abstract JpaRepository<E, UUID> getRepository();
protected void setSearchText(E entity) {
}
@Override
@Transactional
public D save(TenantId tenantId, D domain) {
return save(tenantId, domain, null);
}
@Override
@Transactional
public D saveAndFlush(TenantId tenantId, D domain) {
D d = save(tenantId, domain);
getRepository().flush();
return d;
}
protected D save(TenantId tenantId, D domain, Consumer<E> preSaveAction) {
E entity;
try {
entity = getEntityClass().getConstructor(domain.getClass()).newInstance(domain);
@ -83,7 +60,6 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D extends HasId<?>
log.error("Can't create entity for domain object {}", domain, e);
throw new IllegalArgumentException("Can't create entity for domain object {" + domain + "}", e);
}
setSearchText(entity);
log.debug("Saving entity {}", entity);
boolean isNew = entity.getUuid() == null;
if (isNew) {
@ -92,17 +68,9 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D extends HasId<?>
entity.setCreatedTime(Uuids.unixTimestamp(uuid));
}
if (preSaveAction != null) {
preSaveAction.accept(entity);
}
if (TransactionSynchronizationManager.isActualTransactionActive()) {
return doSave(entity, isNew);
} else {
return transactionTemplate.execute(status -> doSave(entity, isNew));
if (isPartitioned()) {
createPartition(entity);
}
}
private D doSave(E entity, boolean isNew) {
if (isNew) {
entityManager.persist(entity);
} else {
@ -111,6 +79,14 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D extends HasId<?>
return DaoUtil.getData(entity);
}
@Override
@Transactional
public D saveAndFlush(TenantId tenantId, D domain) {
D d = save(tenantId, domain);
getRepository().flush();
return d;
}
@Override
public D findById(TenantId tenantId, UUID key) {
log.debug("Get entity by key {}", key);
@ -155,4 +131,12 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D extends HasId<?>
List<E> entities = Lists.newArrayList(getRepository().findAll());
return DaoUtil.convertDataList(entities);
}
public boolean isPartitioned() {
return false;
}
public void createPartition(E entity) {
}
}

27
dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmCommentDao.java

@ -22,10 +22,8 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.alarm.AlarmComment;
import org.thingsboard.server.common.data.alarm.AlarmCommentInfo;
import org.thingsboard.server.common.data.id.AlarmCommentId;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
@ -54,21 +52,6 @@ public class JpaAlarmCommentDao extends JpaAbstractDao<AlarmCommentEntity, Alarm
@Autowired
private AlarmCommentRepository alarmCommentRepository;
@Transactional
@Override
public AlarmComment createAlarmComment(TenantId tenantId, AlarmComment alarmComment){
log.trace("Saving entity {}", alarmComment);
return save(tenantId, alarmComment, entity -> {
partitioningRepository.createPartitionIfNotExists(ALARM_COMMENT_TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
});
}
@Override
public void deleteAlarmComment(TenantId tenantId, AlarmCommentId alarmCommentId){
log.trace("Try to delete entity alarm comment by id using [{}]", alarmCommentId);
alarmCommentRepository.deleteById(alarmCommentId.getId());
}
@Override
public PageData<AlarmCommentInfo> findAlarmComments(TenantId tenantId, AlarmId id, PageLink pageLink){
log.trace("Try to find alarm comments by alarm id using [{}]", id);
@ -88,6 +71,16 @@ public class JpaAlarmCommentDao extends JpaAbstractDao<AlarmCommentEntity, Alarm
return findByIdAsync(tenantId, key);
}
@Override
public boolean isPartitioned() {
return true;
}
@Override
public void createPartition(AlarmCommentEntity entity) {
partitioningRepository.createPartitionIfNotExists(ALARM_COMMENT_TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
@Override
protected Class<AlarmCommentEntity> getEntityClass() {
return AlarmCommentEntity.class;

24
dao/src/main/java/org/thingsboard/server/dao/sql/audit/JpaAuditLogDao.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.dao.sql.audit;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
@ -26,7 +25,6 @@ import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.audit.AuditLog;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
@ -69,18 +67,6 @@ public class JpaAuditLogDao extends JpaAbstractDao<AuditLogEntity, AuditLog> imp
return auditLogRepository;
}
@Override
public ListenableFuture<AuditLog> saveByTenantId(AuditLog auditLog) {
return service.submit(() -> save(auditLog.getTenantId(), auditLog));
}
@Override
public AuditLog save(TenantId tenantId, AuditLog auditLog) {
return save(tenantId, auditLog, entity -> {
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
});
}
@Override
public PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink) {
return DaoUtil.toPageData(
@ -172,4 +158,14 @@ public class JpaAuditLogDao extends JpaAbstractDao<AuditLogEntity, AuditLog> imp
jdbcTemplate.update("CALL migrate_audit_logs(?, ?, ?)", startTime, endTime, partitionSizeInMs);
}
@Override
public boolean isPartitioned() {
return true;
}
@Override
public void createPartition(AuditLogEntity entity) {
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
}

17
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java

@ -47,7 +47,6 @@ import javax.annotation.PreDestroy;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
@ -151,8 +150,9 @@ public class JpaBaseEdgeEventDao extends JpaAbstractDao<EdgeEventEntity, EdgeEve
if (StringUtils.isEmpty(edgeEvent.getUid())) {
edgeEvent.setUid(edgeEvent.getId().toString());
}
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, edgeEvent.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
return save(new EdgeEventEntity(edgeEvent));
EdgeEventEntity entity = new EdgeEventEntity(edgeEvent);
createPartition(entity);
return save(entity);
}
private ListenableFuture<Void> save(EdgeEventEntity entity) {
@ -227,4 +227,15 @@ public class JpaBaseEdgeEventDao extends JpaAbstractDao<EdgeEventEntity, EdgeEve
private void callMigrationFunction(long startTime, long endTime, long partitionSIzeInMs) {
jdbcTemplate.update("CALL migrate_edge_event(?, ?, ?)", startTime, endTime, partitionSIzeInMs);
}
@Override
public boolean isPartitioned() {
return true;
}
@Override
public void createPartition(EdgeEventEntity entity) {
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
}

21
dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java

@ -19,7 +19,6 @@ import lombok.RequiredArgsConstructor;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.NotificationId;
import org.thingsboard.server.common.data.id.NotificationRequestId;
@ -51,15 +50,6 @@ public class JpaNotificationDao extends JpaAbstractDao<NotificationEntity, Notif
@Value("${sql.notifications.partition_size:168}")
private int partitionSizeInHours;
@Transactional
@Override
public Notification save(TenantId tenantId, Notification notification) {
return save(tenantId, notification, entity -> {
partitioningRepository.createPartitionIfNotExists(ModelConstants.NOTIFICATION_TABLE_NAME,
entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
});
}
@Override
public PageData<Notification> findUnreadByRecipientIdAndPageLink(TenantId tenantId, UserId recipientId, PageLink pageLink) {
return DaoUtil.toPageData(notificationRepository.findByRecipientIdAndStatusNot(recipientId.getId(), NotificationStatus.READ,
@ -110,6 +100,17 @@ public class JpaNotificationDao extends JpaAbstractDao<NotificationEntity, Notif
notificationRepository.deleteByRecipientId(recipientId.getId());
}
@Override
public boolean isPartitioned() {
return true;
}
@Override
public void createPartition(NotificationEntity entity) {
partitioningRepository.createPartitionIfNotExists(ModelConstants.NOTIFICATION_TABLE_NAME,
entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
@Override
protected Class<NotificationEntity> getEntityClass() {
return NotificationEntity.class;

2
dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageInfoDao.java

@ -19,6 +19,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.OtaPackageId;
@ -58,6 +59,7 @@ public class JpaOtaPackageInfoDao extends JpaAbstractDao<OtaPackageInfoEntity, O
return DaoUtil.getData(otaPackageInfoRepository.findOtaPackageInfoById(id));
}
@Transactional
@Override
public OtaPackageInfo save(TenantId tenantId, OtaPackageInfo otaPackageInfo) {
OtaPackageInfo savedOtaPackage = super.save(tenantId, otaPackageInfo);

2
dao/src/test/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmCommentDaoTest.java

@ -85,6 +85,6 @@ public class JpaAlarmCommentDaoTest extends AbstractJpaDaoTest {
alarmComment.setUserId(new UserId(userId));
alarmComment.setType(type);
alarmComment.setComment(JacksonUtil.newObjectNode().put("text", RandomStringUtils.randomAlphanumeric(10)));
alarmCommentDao.createAlarmComment(TenantId.fromUUID(UUID.randomUUID()), alarmComment);
alarmCommentDao.save(TenantId.fromUUID(UUID.randomUUID()), alarmComment);
}
}

Loading…
Cancel
Save