Browse Source

Tasks partitioning strategies; jobs filter; fixes

pull/13285/head
ViacheslavKlimov 1 year ago
parent
commit
f5e816923d
  1. 14
      application/src/main/java/org/thingsboard/server/controller/JobController.java
  2. 25
      application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
  3. 20
      application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java
  4. 5
      application/src/main/resources/thingsboard.yml
  5. 7
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  6. 32
      application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java
  7. 43
      application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java
  8. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java
  9. 2
      common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java
  10. 2
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java
  11. 30
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobFilter.java
  12. 3
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java
  13. 8
      common/data/src/main/java/org/thingsboard/server/common/data/job/task/DummyTask.java
  14. 6
      common/data/src/main/java/org/thingsboard/server/common/data/job/task/Task.java
  15. 2
      common/data/src/main/java/org/thingsboard/server/common/data/job/task/TaskResult.java
  16. 4
      common/proto/src/main/proto/queue.proto
  17. 21
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  18. 2
      common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java
  19. 5
      dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java
  20. 3
      dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java
  21. 28
      dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java
  22. 12
      dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java
  23. 2
      dao/src/main/resources/sql/schema-entities-idx.sql

14
application/src/main/java/org/thingsboard/server/controller/JobController.java

@ -28,12 +28,16 @@ import org.springframework.web.bind.annotation.RestController;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobFilter;
import org.thingsboard.server.common.data.job.JobStatus;
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.job.JobService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.job.JobManager;
import java.util.List;
import java.util.UUID;
import static org.thingsboard.server.controller.ControllerConstants.PAGE_NUMBER_DESCRIPTION;
@ -69,10 +73,16 @@ public class JobController extends BaseController {
@Parameter(description = SORT_PROPERTY_DESCRIPTION)
@RequestParam(required = false) String sortProperty,
@Parameter(description = SORT_ORDER_DESCRIPTION)
@RequestParam(required = false) String sortOrder) throws ThingsboardException {
@RequestParam(required = false) String sortOrder,
@RequestParam(required = false) List<JobType> types,
@RequestParam(required = false) List<JobStatus> statuses) throws ThingsboardException {
// todo check permissions
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder);
return jobService.findJobsByTenantId(getTenantId(), pageLink);
JobFilter filter = JobFilter.builder()
.types(types)
.statuses(statuses)
.build();
return jobService.findJobsByFilter(getTenantId(), filter, pageLink);
}
@PostMapping("/job/{id}/cancel")

25
application/src/main/java/org/thingsboard/server/controller/TelemetryController.java

@ -64,11 +64,9 @@ import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.AggregationParams;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseDeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.DataType;
@ -78,7 +76,6 @@ import org.thingsboard.server.common.data.kv.IntervalType;
import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
@ -317,10 +314,10 @@ public class TelemetryController extends BaseController {
@Parameter(description = "A string value representing the timezone that will be used to calculate exact timestamps for 'WEEK', 'WEEK_ISO', 'MONTH' and 'QUARTER' interval types.")
@RequestParam(name = "timeZone", required = false) String timeZone,
@Parameter(description = "An integer value that represents a max number of time series data points to fetch." +
" This parameter is used only in the case if 'agg' parameter is set to 'NONE'.", schema = @Schema(defaultValue = "100"))
" This parameter is used only in the case if 'agg' parameter is set to 'NONE'.", schema = @Schema(defaultValue = "100"))
@RequestParam(name = "limit", defaultValue = "100") Integer limit,
@Parameter(description = "A string value representing the aggregation function. " +
"If the interval is not specified, 'agg' parameter will use 'NONE' value.",
"If the interval is not specified, 'agg' parameter will use 'NONE' value.",
schema = @Schema(allowableValues = {"MIN", "MAX", "AVG", "SUM", "COUNT", "NONE"}))
@RequestParam(name = "agg", defaultValue = "NONE") String aggStr,
@Parameter(description = SORT_ORDER_DESCRIPTION, schema = @Schema(allowableValues = {"ASC", "DESC"}))
@ -340,12 +337,12 @@ public class TelemetryController extends BaseController {
+ TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH)
@ApiResponses(value = {
@ApiResponse(responseCode = "200", description = SAVE_ATTIRIBUTES_STATUS_OK +
"Platform creates an audit log event about device attributes updates with action type 'ATTRIBUTES_UPDATED', " +
"and also sends event msg to the rule engine with msg type 'ATTRIBUTES_UPDATED'."),
"Platform creates an audit log event about device attributes updates with action type 'ATTRIBUTES_UPDATED', " +
"and also sends event msg to the rule engine with msg type 'ATTRIBUTES_UPDATED'."),
@ApiResponse(responseCode = "400", description = SAVE_ATTIRIBUTES_STATUS_BAD_REQUEST),
@ApiResponse(responseCode = "401", description = "User is not authorized to save device attributes for selected device. Most likely, User belongs to different Customer or Tenant."),
@ApiResponse(responseCode = "500", description = "The exception was thrown during processing the request. " +
"Platform creates an audit log event about device attributes updates with action type 'ATTRIBUTES_UPDATED' that includes an error stacktrace."),
"Platform creates an audit log event about device attributes updates with action type 'ATTRIBUTES_UPDATED' that includes an error stacktrace."),
})
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/{deviceId}/{scope}", method = RequestMethod.POST)
@ -463,11 +460,11 @@ public class TelemetryController extends BaseController {
TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH)
@ApiResponses(value = {
@ApiResponse(responseCode = "200", description = "Time series for the selected keys in the request was removed. " +
"Platform creates an audit log event about entity time series removal with action type 'TIMESERIES_DELETED'."),
"Platform creates an audit log event about entity time series removal with action type 'TIMESERIES_DELETED'."),
@ApiResponse(responseCode = "400", description = "Platform returns a bad request in case if keys list is empty or start and end timestamp values is empty when deleteAllDataForKeys is set to false."),
@ApiResponse(responseCode = "401", description = "User is not authorized to delete entity time series for selected entity. Most likely, User belongs to different Customer or Tenant."),
@ApiResponse(responseCode = "500", description = "The exception was thrown during processing the request. " +
"Platform creates an audit log event about entity time series removal with action type 'TIMESERIES_DELETED' that includes an error stacktrace."),
"Platform creates an audit log event about entity time series removal with action type 'TIMESERIES_DELETED' that includes an error stacktrace."),
})
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/{entityType}/{entityId}/timeseries/delete", method = RequestMethod.DELETE)
@ -544,11 +541,11 @@ public class TelemetryController extends BaseController {
"Referencing a non-existing Device Id will cause an error" + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH)
@ApiResponses(value = {
@ApiResponse(responseCode = "200", description = "Device attributes was removed for the selected keys in the request. " +
"Platform creates an audit log event about device attributes removal with action type 'ATTRIBUTES_DELETED'."),
"Platform creates an audit log event about device attributes removal with action type 'ATTRIBUTES_DELETED'."),
@ApiResponse(responseCode = "400", description = "Platform returns a bad request in case if keys or scope are not specified."),
@ApiResponse(responseCode = "401", description = "User is not authorized to delete device attributes for selected entity. Most likely, User belongs to different Customer or Tenant."),
@ApiResponse(responseCode = "500", description = "The exception was thrown during processing the request. " +
"Platform creates an audit log event about device attributes removal with action type 'ATTRIBUTES_DELETED' that includes an error stacktrace."),
"Platform creates an audit log event about device attributes removal with action type 'ATTRIBUTES_DELETED' that includes an error stacktrace."),
})
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/{deviceId}/{scope}", method = RequestMethod.DELETE)
@ -566,11 +563,11 @@ public class TelemetryController extends BaseController {
INVALID_ENTITY_ID_OR_ENTITY_TYPE_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH)
@ApiResponses(value = {
@ApiResponse(responseCode = "200", description = "Entity attributes was removed for the selected keys in the request. " +
"Platform creates an audit log event about entity attributes removal with action type 'ATTRIBUTES_DELETED'."),
"Platform creates an audit log event about entity attributes removal with action type 'ATTRIBUTES_DELETED'."),
@ApiResponse(responseCode = "400", description = "Platform returns a bad request in case if keys or scope are not specified."),
@ApiResponse(responseCode = "401", description = "User is not authorized to delete entity attributes for selected entity. Most likely, User belongs to different Customer or Tenant."),
@ApiResponse(responseCode = "500", description = "The exception was thrown during processing the request. " +
"Platform creates an audit log event about entity attributes removal with action type 'ATTRIBUTES_DELETED' that includes an error stacktrace."),
"Platform creates an audit log event about entity attributes removal with action type 'ATTRIBUTES_DELETED' that includes an error stacktrace."),
})
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/{entityType}/{entityId}/{scope}", method = RequestMethod.DELETE)

20
application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java

@ -24,6 +24,7 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.rule.engine.api.NotificationCenter;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job;
@ -79,6 +80,8 @@ public class DefaultJobManager implements JobManager {
private final ExecutorService executor;
private final ExecutorService consumerExecutor;
@Value("${queue.tasks.partitioning_strategy:tenant}")
private String tasksPartitioningStrategy;
@Value("${queue.tasks.stats.processing_interval_ms:1000}")
private int statsProcessingInterval;
@ -148,7 +151,7 @@ public class DefaultJobManager implements JobManager {
JobProcessor processor = getJobProcessor(job.getType());
List<TaskResult> toReprocess = job.getConfiguration().getToReprocess();
if (toReprocess == null) {
int tasksCount = processor.process(job, this::submitTask); // todo: think about stopping tb - while tasks are being submitted
int tasksCount = processor.process(job, this::submitTask);
log.info("[{}][{}][{}] Submitted {} tasks", tenantId, jobId, job.getType(), tasksCount);
jobStatsService.reportAllTasksSubmitted(tenantId, jobId, tasksCount);
} else {
@ -197,17 +200,24 @@ public class DefaultJobManager implements JobManager {
}
private void submitTask(Task<?> task) {
log.info("[{}][{}] Submitting task: {}", task.getTenantId(), task.getJobId(), task);
log.debug("[{}][{}] Submitting task: {}", task.getTenantId(), task.getJobId(), task);
TaskProto taskProto = TaskProto.newBuilder()
.setValue(JacksonUtil.toString(task))
.build();
TbQueueProducer<TbProtoQueueMsg<TaskProto>> producer = taskProducers.get(task.getJobType());
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TASK_PROCESSOR, task.getJobType().name(), task.getTenantId(), task.getTenantId()); // one job at a time for a given tenant
EntityId entityId = null;
if (tasksPartitioningStrategy.equals("entity")) {
entityId = task.getEntityId();
}
if (entityId == null) {
entityId = task.getTenantId();
}
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TASK_PROCESSOR, task.getJobType().name(), task.getTenantId(), entityId);
producer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), taskProto), new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
log.trace("Submitted task: {}", task);
log.trace("Submitted task to {}: {}", tpi, taskProto);
}
@Override
@ -247,7 +257,7 @@ public class DefaultJobManager implements JobManager {
});
consumer.commit();
Thread.sleep(statsProcessingInterval); // todo: test with bigger interval
Thread.sleep(statsProcessingInterval);
}
private void sendJobFinishedNotification(Job job) {

5
application/src/main/resources/thingsboard.yml

@ -1900,6 +1900,11 @@ queue:
tasks:
# Partitions count for tasks queues
partitions: "${TB_QUEUE_TASKS_PARTITIONS:12}"
# Custom partitions count for tasks queues per type. Format: 'TYPE1:24;TYPE2:36', e.g. 'CF_REPROCESSING:24;TENANT_EXPORT:6'
partitions_per_type: "${TB_QUEUE_TASKS_PARTITIONS_PER_TYPE:}"
# Tasks partitioning strategy: 'tenant' or 'entity'. By default, using 'tenant' - tasks of a specific tenant are processed in the same partition.
# In a single-tenant environment, use 'entity' strategy to distribute the tasks among multiple partitions.
partitioning_strategy: "${TB_QUEUE_TASKS_PARTITIONING_STRATEGY:tenant}"
stats:
# Interval in milliseconds to process job stats
processing_interval_ms: "${TB_QUEUE_TASKS_STATS_PROCESSING_INTERVAL_MS:1000}"

7
application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java

@ -109,6 +109,7 @@ import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.notification.Notification;
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod;
import org.thingsboard.server.common.data.notification.NotificationType;
@ -168,6 +169,7 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.stream.Collectors;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
@ -1270,6 +1272,11 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
return doGetTypedWithPageLink("/api/jobs?", new TypeReference<PageData<Job>>() {}, new PageLink(100, 0, null, new SortOrder("createdTime", SortOrder.Direction.DESC))).getData();
}
protected List<Job> findJobs(JobType... types) throws Exception {
return doGetTypedWithPageLink("/api/jobs?types=" + Arrays.stream(types).map(Enum::name).collect(Collectors.joining(",")) + "&",
new TypeReference<PageData<Job>>() {}, new PageLink(100, 0, null, new SortOrder("createdTime", SortOrder.Direction.DESC))).getData();
}
protected void cancelJob(JobId jobId) throws Exception {
doPost("/api/job/" + jobId + "/cancel").andExpect(status().isOk());
}

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

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.service.job;
import org.assertj.core.api.Assertions;
import org.assertj.core.api.ThrowingConsumer;
import org.junit.After;
import org.junit.Before;
@ -34,13 +33,12 @@ import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.task.DummyTaskResult;
import org.thingsboard.server.common.data.job.task.DummyTaskResult.DummyTaskFailure;
import org.thingsboard.server.common.data.notification.Notification;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.job.JobService;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.queue.task.JobStatsService;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@ -61,9 +59,6 @@ public class JobManagerTest extends AbstractControllerTest {
@Autowired
private JobManager jobManager;
@Autowired
private JobService jobService;
@SpyBean
private TestTaskProcessor taskProcessor;
@ -138,8 +133,9 @@ public class JobManagerTest extends AbstractControllerTest {
assertThat(jobResult.getSuccessfulCount()).isEqualTo(successfulTasks);
assertThat(jobResult.getFailedCount()).isEqualTo(failedTasks);
assertThat(jobResult.getTotalCount()).isEqualTo(successfulTasks + failedTasks);
assertThat(((DummyTaskResult) jobResult.getResults().get(0)).getFailure().getError()).isEqualTo("error3"); // last error
assertThat(((DummyTaskResult) jobResult.getResults().get(1)).getFailure().getError()).isEqualTo("error3"); // last error
assertThat(getFailures(jobResult)).hasSize(2).allSatisfy(failure -> {
assertThat(failure.getError()).isEqualTo("error3"); // last error
});
assertThat(jobResult.getCompletedCount()).isEqualTo(jobResult.getTotalCount());
});
@ -254,7 +250,7 @@ public class JobManagerTest extends AbstractControllerTest {
Thread.sleep(3000);
verify(jobStatsService, never()).reportTaskResult(any(), any(), any());
Assertions.assertThat(jobService.findJobsByTenantId(tenantId, new PageLink(100, 0)).getData()).isEmpty();
assertThat(findJobs()).isEmpty();
}
@Test
@ -276,7 +272,7 @@ public class JobManagerTest extends AbstractControllerTest {
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {
List<Job> jobs = findJobs();
List<Job> jobs = findJobs(JobType.DUMMY);
assertThat(jobs).hasSize(jobsCount);
Job firstJob = jobs.get(2); // ordered by createdTime descending
assertThat(firstJob.getStatus()).isEqualTo(JobStatus.RUNNING);
@ -391,8 +387,9 @@ public class JobManagerTest extends AbstractControllerTest {
assertThat(jobResult.getSuccessfulCount()).isEqualTo(successfulTasks);
assertThat(jobResult.getFailedCount()).isEqualTo(failedTasks);
List<DummyTaskFailure> failures = getFailures(jobResult);
for (int i = 0, taskNumber = successfulTasks + 1; taskNumber <= totalTasksCount; i++, taskNumber++) {
DummyTaskFailure failure = ((DummyTaskResult) jobResult.getResults().get(i)).getFailure();
DummyTaskFailure failure = failures.get(i);
assertThat(failure.getNumber()).isEqualTo(taskNumber);
assertThat(failure.getError()).isEqualTo("error");
}
@ -438,8 +435,9 @@ public class JobManagerTest extends AbstractControllerTest {
assertThat(jobResult.getFailedCount()).isEqualTo(failedTasks + permanentlyFailedTasks);
assertThat(jobResult.getTotalCount()).isEqualTo(totalTasksCount);
List<DummyTaskFailure> failures = getFailures(jobResult);
for (int i = 0, taskNumber = successfulTasks + 1; taskNumber <= totalTasksCount; i++, taskNumber++) {
DummyTaskFailure failure = ((DummyTaskResult) jobResult.getResults().get(i)).getFailure();
DummyTaskFailure failure = failures.get(i);
assertThat(failure.getNumber()).isEqualTo(taskNumber);
assertThat(failure.getError()).isEqualTo("error");
}
@ -455,8 +453,9 @@ public class JobManagerTest extends AbstractControllerTest {
assertThat(jobResult.getFailedCount()).isEqualTo(permanentlyFailedTasks);
assertThat(jobResult.getTotalCount()).isEqualTo(totalTasksCount);
List<DummyTaskFailure> failures = getFailures(jobResult);
for (int i = 0, taskNumber = successfulTasks + failedTasks + 1; taskNumber <= totalTasksCount; i++, taskNumber++) {
DummyTaskFailure failure = ((DummyTaskResult) jobResult.getResults().get(i)).getFailure();
DummyTaskFailure failure = failures.get(i);
assertThat(failure.getNumber()).isEqualTo(taskNumber);
assertThat(failure.getError()).isEqualTo("error");
assertThat(failure.isFailAlways()).isTrue();
@ -474,6 +473,11 @@ public class JobManagerTest extends AbstractControllerTest {
});
}
// todo: job with zero tasks
private List<DummyTaskFailure> getFailures(JobResult jobResult) {
return jobResult.getResults().stream()
.map(taskResult -> ((DummyTaskResult) taskResult).getFailure())
.sorted(Comparator.comparingInt(DummyTaskFailure::getNumber))
.toList();
}
}

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

@ -0,0 +1,43 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.job;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
@TestPropertySource(properties = {
"queue.tasks.stats.processing_interval_ms=0",
"queue.tasks.partitioning_strategy=entity",
"queue.tasks.partitions_per_type=DUMMY:100;DUMMY:50"
})
public class JobManagerTest_EntityPartitioningStrategy extends JobManagerTest {
/*
* Some tests are overridden because they are based on
* tenant partitioning strategy (subsequent tasks processing within a tenant)
* */
@Override
public void testCancelJob_simulateTaskProcessorRestart() throws Exception {
}
@Override
public void testGeneralJobError() {
}
}

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

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.job;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobFilter;
import org.thingsboard.server.common.data.job.JobStats;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
@ -35,7 +36,7 @@ public interface JobService extends EntityDaoService {
void processStats(TenantId tenantId, JobId jobId, JobStats jobStats);
PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink);
PageData<Job> findJobsByFilter(TenantId tenantId, JobFilter filter, PageLink pageLink);
Job findLatestJobByKey(TenantId tenantId, String key);

2
common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java

@ -29,7 +29,7 @@ public class JobId extends UUIDBased implements EntityId {
super(id);
}
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "string", example = "TASK", allowableValues = "TASK")
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "string", example = "JOB", allowableValues = "JOB")
@Override
public EntityType getEntityType() {
return EntityType.JOB;

2
common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.data.job;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
@ -35,6 +36,7 @@ public abstract class JobConfiguration implements Serializable {
private List<TaskResult> toReprocess;
@JsonIgnore
public abstract JobType getType();
}

30
common/data/src/main/java/org/thingsboard/server/common/data/job/JobFilter.java

@ -0,0 +1,30 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.data.job;
import lombok.Builder;
import lombok.Data;
import java.util.List;
@Data
@Builder
public class JobFilter {
private final List<JobType> types;
private final List<JobStatus> statuses;
}

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

@ -58,7 +58,7 @@ public abstract class JobResult implements Serializable {
discardedCount++;
} else {
failedCount++;
if (results.size() < 1000) { // preserving only first 1000 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);
}
}
@ -67,6 +67,7 @@ public abstract class JobResult implements Serializable {
@JsonIgnore
public abstract String getDescription();
@JsonIgnore
public abstract JobType getJobType();
}

8
common/data/src/main/java/org/thingsboard/server/common/data/job/task/DummyTask.java

@ -20,9 +20,12 @@ import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import lombok.ToString;
import lombok.experimental.SuperBuilder;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.job.JobType;
import java.util.List;
import java.util.UUID;
@Data
@NoArgsConstructor
@ -46,6 +49,11 @@ public class DummyTask extends Task<DummyTaskResult> {
return DummyTaskResult.discarded();
}
@Override
public EntityId getEntityId() {
return new DeviceId(UUID.randomUUID());
}
@Override
public JobType getJobType() {
return JobType.DUMMY;

6
common/data/src/main/java/org/thingsboard/server/common/data/job/task/Task.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.data.job.task;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
@ -22,6 +23,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.experimental.SuperBuilder;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.JobType;
@ -49,6 +51,10 @@ public abstract class Task<R extends TaskResult> {
public abstract R toDiscarded();
@JsonIgnore
public abstract EntityId getEntityId();
@JsonIgnore
public abstract JobType getJobType();
}

2
common/data/src/main/java/org/thingsboard/server/common/data/job/task/TaskResult.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.data.job.task;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
@ -39,6 +40,7 @@ public abstract class TaskResult {
private boolean success;
private boolean discarded;
@JsonIgnore
public abstract JobType getJobType();
}

4
common/proto/src/main/proto/queue.proto

@ -1854,7 +1854,7 @@ message EdqsResponseMsg {
}
message TaskProto {
string value = 1; // fixme: TMP, make more efficient
string value = 1;
}
message JobStatsMsg {
@ -1867,5 +1867,5 @@ message JobStatsMsg {
}
message TaskResultProto {
string value = 1; // fixme: TMP, make more efficient
string value = 1;
}

21
common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

@ -40,12 +40,14 @@ import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.PropertyUtils;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.EnumMap;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@ -87,7 +89,9 @@ public class HashPartitionService implements PartitionService {
@Value("${queue.edqs.partitions:12}")
private Integer edqsPartitions;
@Value("${queue.tasks.partitions:12}")
private Integer tasksPartitions;
private Integer defaultTasksPartitions;
@Value("${queue.tasks.partitions_per_type:}")
private String tasksPartitionsPerType;
@Value("${queue.partitions.hash_function_name:murmur3_128}")
private String hashFunctionName;
@ -135,11 +139,18 @@ public class HashPartitionService implements PartitionService {
partitionSizesMap.put(edqsKey, edqsPartitions);
partitionTopicsMap.put(edqsKey, "edqs"); // placeholder, not used
for (JobType jobType : JobType.values()) {
QueueKey queueKey = new QueueKey(ServiceType.TASK_PROCESSOR, jobType.name());
partitionSizesMap.put(queueKey, tasksPartitions);
partitionTopicsMap.put(queueKey, jobType.getTasksTopic());
Map<JobType, Integer> tasksPartitions = new EnumMap<>(JobType.class);
PropertyUtils.getProps(tasksPartitionsPerType).forEach((type, partitions) -> {
tasksPartitions.put(JobType.valueOf(type), Integer.parseInt(partitions));
});
for (JobType type : JobType.values()) {
tasksPartitions.putIfAbsent(type, defaultTasksPartitions);
}
tasksPartitions.forEach((type, partitions) -> {
QueueKey queueKey = new QueueKey(ServiceType.TASK_PROCESSOR, type.name());
partitionSizesMap.put(queueKey, partitions);
partitionTopicsMap.put(queueKey, type.getTasksTopic());
});
}
@AfterStartUp(order = AfterStartUp.QUEUE_INFO_INITIALIZATION)

2
common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java

@ -63,7 +63,7 @@ public abstract class TaskProcessor<T extends Task<R>, R extends TaskResult> {
private QueueKey queueKey;
private MainQueueConsumerManager<TbProtoQueueMsg<TaskProto>, QueueConfig> taskConsumer;
private final ExecutorService taskExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-processor"));
private final ExecutorService taskExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-processor"));
private final SetCache<UUID> discardedJobs = new SetCache<>(TimeUnit.MINUTES.toMillis(60));
private final SetCache<UUID> deletedTenants = new SetCache<>(TimeUnit.MINUTES.toMillis(60));

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

@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobFilter;
import org.thingsboard.server.common.data.job.JobResult;
import org.thingsboard.server.common.data.job.JobStats;
import org.thingsboard.server.common.data.job.JobStatus;
@ -174,8 +175,8 @@ public class DefaultJobService extends AbstractEntityService implements JobServi
}
@Override
public PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink) {
return jobDao.findByTenantId(tenantId, pageLink);
public PageData<Job> findJobsByFilter(TenantId tenantId, JobFilter filter, PageLink pageLink) {
return jobDao.findByTenantIdAndFilter(tenantId, filter, pageLink);
}
@Override

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

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.job;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobFilter;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.page.PageData;
@ -26,7 +27,7 @@ import org.thingsboard.server.dao.Dao;
public interface JobDao extends Dao<Job> {
PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink);
PageData<Job> findByTenantIdAndFilter(TenantId tenantId, JobFilter filter, PageLink pageLink);
Job findByIdForUpdate(TenantId tenantId, JobId jobId);

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

@ -15,12 +15,9 @@
*/
package org.thingsboard.server.dao.sql.job;
import jakarta.persistence.LockModeType;
import org.springframework.data.domain.Limit;
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;
@ -37,28 +34,29 @@ 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)")
Page<JobEntity> findByTenantIdAndSearchText(@Param("tenantId") UUID tenantId,
@Param("searchText") String searchText,
Pageable pageable);
"AND (:types IS NULL OR j.type IN (:types)) AND (:statuses IS NULL OR j.status IN (:statuses)) " +
"AND (:searchText IS NULL OR ilike(j.key, concat('%', :searchText, '%')) = true " +
"OR ilike(j.description, concat('%', :searchText, '%')) = true)")
Page<JobEntity> findByTenantIdAndTypesAndStatusesAndSearchText(@Param("tenantId") UUID tenantId,
@Param("types") List<JobType> types,
@Param("statuses") List<JobStatus> statuses,
@Param("searchText") String searchText,
Pageable pageable);
@Lock(LockModeType.PESSIMISTIC_WRITE) // SELECT FOR UPDATE
@Query("SELECT j FROM JobEntity j WHERE j.id = :id")
@Query(value = "SELECT * FROM job j WHERE j.id = :id FOR UPDATE", nativeQuery = true)
JobEntity findByIdForUpdate(UUID id);
@Query("SELECT j FROM JobEntity j WHERE j.tenantId = :tenantId AND j.key = :key " +
"ORDER BY j.createdTime DESC")
"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")
JobEntity findOldestByTenantIdAndTypeAndStatusForUpdate(UUID tenantId, JobType type, JobStatus status, Limit limit);
@Query(value = "SELECT * FROM job j WHERE j.tenant_id = :tenantId AND j.type = :type " +
"AND j.status = :status ORDER BY j.created_time ASC, j.id ASC LIMIT 1 FOR UPDATE", nativeQuery = true)
JobEntity findOldestByTenantIdAndTypeAndStatusForUpdate(UUID tenantId, String type, String status);
@Transactional
@Modifying

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

@ -17,17 +17,18 @@ package org.thingsboard.server.dao.sql.job;
import com.google.common.base.Strings;
import lombok.RequiredArgsConstructor;
import org.springframework.data.domain.Limit;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobFilter;
import org.thingsboard.server.common.data.job.JobStatus;
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.common.data.util.CollectionsUtil;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.job.JobDao;
import org.thingsboard.server.dao.model.sql.JobEntity;
@ -45,8 +46,11 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
private final JobRepository jobRepository;
@Override
public PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink) {
return DaoUtil.toPageData(jobRepository.findByTenantIdAndSearchText(tenantId.getId(), Strings.emptyToNull(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink)));
public PageData<Job> findByTenantIdAndFilter(TenantId tenantId, JobFilter filter, PageLink pageLink) {
return DaoUtil.toPageData(jobRepository.findByTenantIdAndTypesAndStatusesAndSearchText(tenantId.getId(),
CollectionsUtil.isEmpty(filter.getTypes()) ? null : filter.getTypes(),
CollectionsUtil.isEmpty(filter.getStatuses()) ? null : filter.getStatuses(),
Strings.emptyToNull(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink)));
}
@Override
@ -71,7 +75,7 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
@Override
public Job findOldestByTenantIdAndTypeAndStatusForUpdate(TenantId tenantId, JobType type, JobStatus status) {
return DaoUtil.getData(jobRepository.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId.getId(), type, status, Limit.of(1)));
return DaoUtil.getData(jobRepository.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId.getId(), type.name(), status.name()));
}
@Override

2
dao/src/main/resources/sql/schema-entities-idx.sql

@ -129,3 +129,5 @@ CREATE INDEX IF NOT EXISTS idx_resource_etag ON resource(tenant_id, etag);
CREATE INDEX IF NOT EXISTS idx_resource_type_public_resource_key ON resource(resource_type, public_resource_key);
CREATE INDEX IF NOT EXISTS mobile_app_bundle_tenant_id ON mobile_app_bundle(tenant_id);
CREATE INDEX IF NOT EXISTS idx_job_tenant_id ON job(tenant_id);

Loading…
Cancel
Save