diff --git a/application/src/main/java/org/thingsboard/server/controller/JobController.java b/application/src/main/java/org/thingsboard/server/controller/JobController.java index 9b6627e12a..1db84593f5 100644 --- a/application/src/main/java/org/thingsboard/server/controller/JobController.java +++ b/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 types, + @RequestParam(required = false) List 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") diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index d93e973073..bf9713f58e 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/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) diff --git a/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java b/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java index 79312ba608..8c32de4f45 100644 --- a/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java +++ b/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 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> 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) { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 54db919bbf..7cae5db294 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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}" diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 054c2fbe59..ef658fc17a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/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>() {}, new PageLink(100, 0, null, new SortOrder("createdTime", SortOrder.Direction.DESC))).getData(); } + protected List findJobs(JobType... types) throws Exception { + return doGetTypedWithPageLink("/api/jobs?types=" + Arrays.stream(types).map(Enum::name).collect(Collectors.joining(",")) + "&", + new TypeReference>() {}, 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()); } diff --git a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java index 15037470d1..206acc450d 100644 --- a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java +++ b/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 jobs = findJobs(); + List 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 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 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 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 getFailures(JobResult jobResult) { + return jobResult.getResults().stream() + .map(taskResult -> ((DummyTaskResult) taskResult).getFailure()) + .sorted(Comparator.comparingInt(DummyTaskFailure::getNumber)) + .toList(); + } } \ No newline at end of file diff --git a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java new file mode 100644 index 0000000000..983f30d523 --- /dev/null +++ b/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() { + + } + +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java index 3c00b2fe0e..3204044880 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/job/JobService.java +++ b/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 findJobsByTenantId(TenantId tenantId, PageLink pageLink); + PageData findJobsByFilter(TenantId tenantId, JobFilter filter, PageLink pageLink); Job findLatestJobByKey(TenantId tenantId, String key); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java index e6688f0eb0..76678b8b31 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java +++ b/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; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java index 8aed4adbe5..35f44cd4be 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java +++ b/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 toReprocess; + @JsonIgnore public abstract JobType getType(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobFilter.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobFilter.java new file mode 100644 index 0000000000..6cc9a636e8 --- /dev/null +++ b/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 types; + private final List statuses; + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java index 534e3587bb..3af076cd19 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java +++ b/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(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/DummyTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/task/DummyTask.java index d7a37b5175..7e262ed7b8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/DummyTask.java +++ b/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 { return DummyTaskResult.discarded(); } + @Override + public EntityId getEntityId() { + return new DeviceId(UUID.randomUUID()); + } + @Override public JobType getJobType() { return JobType.DUMMY; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/Task.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/task/Task.java index 6bf287b54a..cce32fdad0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/Task.java +++ b/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 { public abstract R toDiscarded(); + @JsonIgnore + public abstract EntityId getEntityId(); + + @JsonIgnore public abstract JobType getJobType(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/TaskResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/task/TaskResult.java index 8c4667e1be..9416cb8f6b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/task/TaskResult.java +++ b/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(); } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 1cc507c5b3..04c41a7f69 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/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; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 3f81199488..d6580d5c83 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/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 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) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java b/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java index 38d3d02b50..fd4367781e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java @@ -63,7 +63,7 @@ public abstract class TaskProcessor, R extends TaskResult> { private QueueKey queueKey; private MainQueueConsumerManager, 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 discardedJobs = new SetCache<>(TimeUnit.MINUTES.toMillis(60)); private final SetCache deletedTenants = new SetCache<>(TimeUnit.MINUTES.toMillis(60)); diff --git a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java b/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java index 73e0c5dff4..680bc81566 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java +++ b/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 findJobsByTenantId(TenantId tenantId, PageLink pageLink) { - return jobDao.findByTenantId(tenantId, pageLink); + public PageData findJobsByFilter(TenantId tenantId, JobFilter filter, PageLink pageLink) { + return jobDao.findByTenantIdAndFilter(tenantId, filter, pageLink); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java b/dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java index afe182d8cd..0c70fd102d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/job/JobDao.java +++ b/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 { - PageData findByTenantId(TenantId tenantId, PageLink pageLink); + PageData findByTenantIdAndFilter(TenantId tenantId, JobFilter filter, PageLink pageLink); Job findByIdForUpdate(TenantId tenantId, JobId jobId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java index 0ecd517f51..72d569c94b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/job/JobRepository.java +++ b/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 { @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 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 findByTenantIdAndTypesAndStatusesAndSearchText(@Param("tenantId") UUID tenantId, + @Param("types") List types, + @Param("statuses") List 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 statuses); boolean existsByTenantIdAndTypeAndStatusIn(UUID tenantId, JobType type, List 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 diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java index 0e5ee4683e..f59898eb73 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/job/JpaJobDao.java +++ b/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 implements JobDao private final JobRepository jobRepository; @Override - public PageData findByTenantId(TenantId tenantId, PageLink pageLink) { - return DaoUtil.toPageData(jobRepository.findByTenantIdAndSearchText(tenantId.getId(), Strings.emptyToNull(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink))); + public PageData 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 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 diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 7f52365e33..ad311f00df 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/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);