Browse Source

Merge pull request #13293 from irynamatveieva/jobs

Added findByKey method to job service
pull/13285/head
Viacheslav Klimov 1 year ago
committed by GitHub
parent
commit
4057e6b42e
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java
  2. 14
      dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java
  3. 6
      dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java
  4. 19
      dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java
  5. 16
      dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java

2
common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java

@ -37,4 +37,6 @@ public interface JobService extends EntityDaoService {
PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink);
Job findLatestJobByKey(TenantId tenantId, String key);
}

14
dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java

@ -29,8 +29,8 @@ import org.thingsboard.server.common.data.job.JobResult;
import org.thingsboard.server.common.data.job.JobStats;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.job.TaskFailure;
import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.entity.AbstractEntityService;
@ -55,7 +55,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
@Transactional
@Override
public Job submitJob(TenantId tenantId, Job job) {
if (jobDao.existsByKeyAndStatusOneOf(job.getKey(), QUEUED, PENDING, RUNNING)) {
if (jobDao.existsByTenantAndKeyAndStatusOneOf(tenantId, job.getKey(), QUEUED, PENDING, RUNNING)) {
throw new IllegalArgumentException("The same job is already queued or running");
}
if (jobDao.existsByTenantIdAndTypeAndStatusOneOf(tenantId, job.getType(), PENDING, RUNNING)) {
@ -187,6 +187,11 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
return jobDao.findByTenantId(tenantId, pageLink);
}
@Override
public Job findLatestJobByKey(TenantId tenantId, String key) {
return jobDao.findLatestByKey(tenantId, key);
}
private Job findForUpdate(TenantId tenantId, JobId jobId) {
return jobDao.findByIdForUpdate(tenantId, jobId);
}
@ -201,6 +206,11 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
jobDao.removeById(tenantId, id.getId());
}
@Override
public void deleteByTenantId(TenantId tenantId) {
jobDao.deleteByTenantId(tenantId);
}
@Override
public EntityType getEntityType() {
return EntityType.JOB;

6
dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java

@ -30,10 +30,14 @@ public interface JobDao extends Dao<Job> {
Job findByIdForUpdate(TenantId tenantId, JobId jobId);
boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses);
Job findLatestByKey(TenantId tenantId, String key);
boolean existsByTenantAndKeyAndStatusOneOf(TenantId tenantId, String key, JobStatus... statuses);
boolean existsByTenantIdAndTypeAndStatusOneOf(TenantId tenantId, JobType type, JobStatus... statuses);
Job findOldestByTenantIdAndTypeAndStatusForUpdate(TenantId tenantId, JobType type, JobStatus status);
void deleteByTenantId(TenantId tenantId);
}

19
dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java

@ -21,9 +21,11 @@ import org.springframework.data.domain.Page;
import org.springframework.data.domain.Pageable;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Lock;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.dao.model.sql.JobEntity;
@ -35,8 +37,8 @@ import java.util.UUID;
public interface JobRepository extends JpaRepository<JobEntity, UUID> {
@Query("SELECT j FROM JobEntity j WHERE j.tenantId = :tenantId " +
"AND (:searchText IS NULL OR ilike(j.key, concat('%', :searchText, '%')) = true " +
"OR ilike(j.description, concat('%', :searchText, '%')) = true)")
"AND (:searchText IS NULL OR ilike(j.key, concat('%', :searchText, '%')) = true " +
"OR ilike(j.description, concat('%', :searchText, '%')) = true)")
Page<JobEntity> findByTenantIdAndSearchText(@Param("tenantId") UUID tenantId,
@Param("searchText") String searchText,
Pageable pageable);
@ -45,13 +47,22 @@ public interface JobRepository extends JpaRepository<JobEntity, UUID> {
@Query("SELECT j FROM JobEntity j WHERE j.id = :id")
JobEntity findByIdForUpdate(UUID id);
boolean existsByKeyAndStatusIn(String key, List<JobStatus> statuses);
@Query("SELECT j FROM JobEntity j WHERE j.tenantId = :tenantId AND j.key = :key " +
"ORDER BY j.createdTime DESC")
JobEntity findLatestByTenantIdAndKey(@Param("tenantId") UUID tenantId, @Param("key") String key);
boolean existsByTenantIdAndKeyAndStatusIn(UUID tenantId, String key, List<JobStatus> statuses);
boolean existsByTenantIdAndTypeAndStatusIn(UUID tenantId, JobType type, List<JobStatus> statuses);
@Lock(LockModeType.PESSIMISTIC_WRITE) // SELECT FOR UPDATE
@Query("SELECT j FROM JobEntity j WHERE j.tenantId = :tenantId AND j.type = :type " +
"AND j.status = :status ORDER BY j.createdTime ASC, j.id ASC")
"AND j.status = :status ORDER BY j.createdTime ASC, j.id ASC")
JobEntity findOldestByTenantIdAndTypeAndStatusForUpdate(UUID tenantId, JobType type, JobStatus status, Limit limit);
@Transactional
@Modifying
@Query("DELETE FROM JobEntity j WHERE j.tenantId = :tenantId")
void deleteByTenantId(UUID tenantId);
}

16
dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java

@ -29,9 +29,9 @@ import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.job.JobDao;
import org.thingsboard.server.dao.model.sql.JobEntity;
import org.thingsboard.server.dao.sql.JpaAbstractDao;
import org.thingsboard.server.dao.job.JobDao;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.Arrays;
@ -55,8 +55,13 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
}
@Override
public boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses) {
return jobRepository.existsByKeyAndStatusIn(key, Arrays.stream(statuses).toList());
public Job findLatestByKey(TenantId tenantId, String key) {
return DaoUtil.getData(jobRepository.findLatestByTenantIdAndKey(tenantId.getId(), key));
}
@Override
public boolean existsByTenantAndKeyAndStatusOneOf(TenantId tenantId, String key, JobStatus... statuses) {
return jobRepository.existsByTenantIdAndKeyAndStatusIn(tenantId.getId(), key, Arrays.stream(statuses).toList());
}
@Override
@ -69,6 +74,11 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
return DaoUtil.getData(jobRepository.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId.getId(), type, status, Limit.of(1)));
}
@Override
public void deleteByTenantId(TenantId tenantId) {
jobRepository.deleteByTenantId(tenantId.getId());
}
@Override
public EntityType getEntityType() {
return EntityType.JOB;

Loading…
Cancel
Save