Browse Source

Merge pull request #14264 from thingsboard/fix/job-cancellation

Task manager: improve task cancellation handling
pull/14267/head
Viacheslav Klimov 11 months ago
committed by GitHub
parent
commit
0c9afa38ff
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java
  2. 9
      application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java
  3. 12
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java
  4. 10
      dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java

2
application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java

@ -98,7 +98,7 @@ public class JobManagerTest extends AbstractControllerTest {
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {
Job job = findJobById(jobId); Job job = findJobById(jobId);
assertThat(job.getStatus()).isEqualTo(JobStatus.RUNNING); assertThat(job.getStatus()).isEqualTo(JobStatus.RUNNING);
assertThat(job.getResult().getSuccessfulCount()).isBetween(1, tasksCount - 1); assertThat(job.getResult().getSuccessfulCount()).isBetween(0, tasksCount - 1);
assertThat(job.getResult().getTotalCount()).isEqualTo(tasksCount); assertThat(job.getResult().getTotalCount()).isEqualTo(tasksCount);
}); });
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {

9
application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java

@ -27,17 +27,16 @@ import org.thingsboard.server.dao.service.DaoSqlTest;
public class JobManagerTest_EntityPartitioningStrategy extends JobManagerTest { public class JobManagerTest_EntityPartitioningStrategy extends JobManagerTest {
/* /*
* Some tests are overridden because they are based on * Some tests are overridden because they are based on
* tenant partitioning strategy (subsequent tasks processing within a tenant) * tenant partitioning strategy (subsequent tasks processing within a tenant)
* */ * */
@Override @Override
public void testCancelJob_simulateTaskProcessorRestart() throws Exception { public void testCancelJob_simulateTaskProcessorRestart() {
} }
@Override @Override
public void testSubmitJob_generalError() { public void testSubmitJob_generalError() {
} }
} }

12
common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java

@ -55,11 +55,17 @@ public abstract class JobResult implements Serializable {
public void processTaskResult(TaskResult taskResult) { public void processTaskResult(TaskResult taskResult) {
if (taskResult.isSuccess()) { if (taskResult.isSuccess()) {
successfulCount++; if (totalCount == null || successfulCount < totalCount) {
successfulCount++;
}
} else if (taskResult.isDiscarded()) { } else if (taskResult.isDiscarded()) {
discardedCount++; if (totalCount == null || discardedCount < totalCount) {
discardedCount++;
}
} else { } else {
failedCount++; if (totalCount == null || failedCount < totalCount) {
failedCount++;
}
if (results.size() < 100) { // preserving only first 100 errors, not reprocessing if there are more failures if (results.size() < 100) { // preserving only first 100 errors, not reprocessing if there are more failures
results.add(taskResult); results.add(taskResult);
} }

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

@ -87,11 +87,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
} }
job.getResult().setCancellationTs(System.currentTimeMillis()); job.getResult().setCancellationTs(System.currentTimeMillis());
JobStatus prevStatus = job.getStatus(); JobStatus prevStatus = job.getStatus();
if (job.getStatus() == QUEUED) { job.setStatus(CANCELLED);
job.setStatus(CANCELLED); // setting cancelled status right away, because we don't expect stats for cancelled tasks
} else if (job.getStatus() == PENDING) {
job.setStatus(RUNNING);
}
saveJob(tenantId, job, true, prevStatus); saveJob(tenantId, job, true, prevStatus);
} }
@ -145,7 +141,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
} }
} }
if (job.getStatus() == RUNNING) { if (job.getStatus().isOneOf(RUNNING, CANCELLED)) {
if (result.getTotalCount() != null && result.getCompletedCount() >= result.getTotalCount()) { if (result.getTotalCount() != null && result.getCompletedCount() >= result.getTotalCount()) {
if (result.getCancellationTs() > 0) { if (result.getCancellationTs() > 0) {
job.setStatus(CANCELLED); job.setStatus(CANCELLED);
@ -193,7 +189,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
private void checkWaitingJobs(TenantId tenantId, JobType jobType) { private void checkWaitingJobs(TenantId tenantId, JobType jobType) {
Job queuedJob = jobDao.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId, jobType, QUEUED); Job queuedJob = jobDao.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId, jobType, QUEUED);
if (queuedJob == null) { if (queuedJob == null || jobDao.existsByTenantIdAndTypeAndStatusOneOf(tenantId, jobType, PENDING, RUNNING)) {
return; return;
} }
queuedJob.setStatus(PENDING); queuedJob.setStatus(PENDING);

Loading…
Cancel
Save