Browse Source

Job stats, REST API, tests

pull/13286/head
ViacheslavKlimov 1 year ago
parent
commit
290fba4819
  1. 73
      application/src/main/java/org/thingsboard/server/controller/JobController.java
  2. 37
      application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java
  3. 64
      application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java
  4. 64
      application/src/main/java/org/thingsboard/server/service/job/DummyJobProcessor.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/job/JobManager.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java
  7. 44
      application/src/main/java/org/thingsboard/server/service/job/task/DummyTaskProcessor.java
  8. 4
      application/src/main/resources/thingsboard.yml
  9. 125
      application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java
  10. 10
      common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java
  11. 2
      common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java
  12. 16
      common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java
  13. 42
      common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobConfiguration.java
  14. 25
      common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobResult.java
  15. 40
      common/data/src/main/java/org/thingsboard/server/common/data/job/DummyTask.java
  16. 12
      common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java
  17. 1
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java
  18. 8
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java
  19. 29
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobStats.java
  20. 3
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java
  21. 15
      common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java
  22. 4
      common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java
  23. 7
      common/proto/src/main/proto/queue.proto
  24. 9
      common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java
  25. 22
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java
  26. 37
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java
  27. 10
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java
  28. 4
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java
  29. 6
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java
  30. 8
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java
  31. 3
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java
  32. 8
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java
  33. 6
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java
  34. 5
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java
  35. 63
      common/queue/src/main/java/org/thingsboard/server/queue/task/JobStatsService.java
  36. 31
      common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java
  37. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  38. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java
  39. 21
      dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java
  40. 16
      dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java
  41. 40
      dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java
  42. 6
      dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java
  43. 1
      dao/src/main/resources/sql/schema-entities.sql

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

@ -0,0 +1,73 @@
/**
* 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.controller;
import io.swagger.v3.oas.annotations.Parameter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
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.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.task.JobService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.UUID;
import static org.thingsboard.server.controller.ControllerConstants.PAGE_NUMBER_DESCRIPTION;
import static org.thingsboard.server.controller.ControllerConstants.PAGE_SIZE_DESCRIPTION;
import static org.thingsboard.server.controller.ControllerConstants.SORT_ORDER_DESCRIPTION;
import static org.thingsboard.server.controller.ControllerConstants.SORT_PROPERTY_DESCRIPTION;
@RestController
@TbCoreComponent
@RequestMapping("/api")
@RequiredArgsConstructor
@Slf4j
public class JobController extends BaseController {
private final JobService jobService;
@GetMapping("/job/{id}")
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
public Job getJobById(@PathVariable UUID id) throws ThingsboardException {
return jobService.findJobById(getTenantId(), new JobId(id));
}
@GetMapping("/jobs")
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
public PageData<Job> getJobs(@Parameter(description = PAGE_SIZE_DESCRIPTION, required = true)
@RequestParam int pageSize,
@Parameter(description = PAGE_NUMBER_DESCRIPTION, required = true)
@RequestParam int page,
@Parameter(description = "Case-insensitive 'substring' filter based on job's description")
@RequestParam(required = false) String textSearch,
@Parameter(description = SORT_PROPERTY_DESCRIPTION)
@RequestParam(required = false) String sortProperty,
@Parameter(description = SORT_ORDER_DESCRIPTION)
@RequestParam(required = false) String sortOrder) throws ThingsboardException {
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder);
return jobService.findJobsByTenantId(getTenantId(), pageLink);
}
}

37
application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java

@ -17,9 +17,11 @@ package org.thingsboard.server.service.job;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.ProfileEntityIdInfo;
import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.id.AssetProfileId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.job.CfReprocessingJobConfiguration; import org.thingsboard.server.common.data.job.CfReprocessingJobConfiguration;
import org.thingsboard.server.common.data.job.CfReprocessingTask; import org.thingsboard.server.common.data.job.CfReprocessingTask;
@ -39,28 +41,34 @@ public class CfReprocessingJobProcessor extends JobProcessor {
private final DeviceService deviceService; private final DeviceService deviceService;
private final AssetService assetService; private final AssetService assetService;
// fixme: multiple jobs with single type
@Transactional
@Override @Override
public void process(Job job, Consumer<Task> taskConsumer) { public int process(Job job, Consumer<Task> taskConsumer) {
CfReprocessingJobConfiguration configuration = job.getConfiguration(); CfReprocessingJobConfiguration configuration = job.getConfiguration();
CalculatedField calculatedField = configuration.getCalculatedField(); CalculatedField calculatedField = configuration.getCalculatedField();
EntityId entityId = calculatedField.getEntityId(); EntityId cfEntityId = calculatedField.getEntityId();
if (entityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) { int tasksCount = 0;
taskConsumer.accept(createTask(job, configuration, entityId)); if (cfEntityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) {
taskConsumer.accept(createTask(job, configuration, cfEntityId));
tasksCount++;
} else { } else {
PageDataIterable<ProfileEntityIdInfo> entities; PageDataIterable<? extends EntityId> entities;
if (entityId.getEntityType() == EntityType.DEVICE_PROFILE) { if (cfEntityId.getEntityType() == EntityType.DEVICE_PROFILE) {
entities = new PageDataIterable<>(pageLink -> deviceService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512); entities = new PageDataIterable<>(pageLink -> deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(job.getTenantId(), (DeviceProfileId) cfEntityId, pageLink), 512);
} else if (entityId.getEntityType() == EntityType.ASSET_PROFILE) { } else if (cfEntityId.getEntityType() == EntityType.ASSET_PROFILE) {
entities = new PageDataIterable<>(pageLink -> assetService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512); entities = new PageDataIterable<>(pageLink -> assetService.findAssetIdsByTenantIdAndAssetProfileId(job.getTenantId(), (AssetProfileId) cfEntityId, pageLink), 512);
} else { } else {
throw new IllegalArgumentException("Unsupported CF entity type " + entityId.getEntityType()); throw new IllegalArgumentException("Unsupported CF entity type " + cfEntityId.getEntityType());
}
for (EntityId entityId : entities) {
taskConsumer.accept(createTask(job, configuration, entityId));
tasksCount++;
} }
entities.forEach(device -> {
taskConsumer.accept(createTask(job, configuration, device.getEntityId()));
});
} }
return tasksCount;
} }
private Task createTask(Job job, CfReprocessingJobConfiguration configuration, EntityId entityId) { private Task createTask(Job job, CfReprocessingJobConfiguration configuration, EntityId entityId) {
@ -68,6 +76,7 @@ public class CfReprocessingJobProcessor extends JobProcessor {
.tenantId(job.getTenantId()) .tenantId(job.getTenantId())
.jobId(job.getId()) .jobId(job.getId())
.key(entityId.getEntityType().getNormalName() + " " + entityId.getId()) .key(entityId.getEntityType().getNormalName() + " " + entityId.getId())
.retries(2) // 3 attempts in total
.calculatedField(configuration.getCalculatedField()) .calculatedField(configuration.getCalculatedField())
.entityId(entityId) .entityId(entityId)
.startTs(configuration.getStartTs()) .startTs(configuration.getStartTs())

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

@ -18,18 +18,20 @@ package org.thingsboard.server.service.job;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
import lombok.SneakyThrows; import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.job.Job; import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobStats;
import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.Task; import org.thingsboard.server.common.data.job.Task;
import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.task.JobService; import org.thingsboard.server.dao.task.JobService;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueMsgMetadata; import org.thingsboard.server.queue.TbQueueMsgMetadata;
@ -37,12 +39,15 @@ import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.task.JobStatsService;
import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.Arrays; import java.util.Arrays;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.function.Function; import java.util.function.Function;
@ -54,23 +59,26 @@ import java.util.stream.Collectors;
public class DefaultJobManager implements JobManager { public class DefaultJobManager implements JobManager {
private final JobService jobService; private final JobService jobService;
private final TbCoreQueueFactory queueFactory; private final JobStatsService jobStatsService;
private final Map<JobType, JobProcessor> jobProcessors; private final Map<JobType, JobProcessor> jobProcessors;
private final Map<JobType, TbQueueProducer<TbProtoQueueMsg<TaskProto>>> taskProducers; private final Map<JobType, TbQueueProducer<TbProtoQueueMsg<TaskProto>>> taskProducers;
private final QueueConsumerManager<TbProtoQueueMsg<TaskResultProto>> taskResultConsumer; private final QueueConsumerManager<TbProtoQueueMsg<JobStatsMsg>> taskResultConsumer;
private final ExecutorService consumerExecutor; private final ExecutorService consumerExecutor;
public DefaultJobManager(JobService jobService, TbCoreQueueFactory queueFactory, List<JobProcessor> jobProcessors) { @Value("${queue.tasks.stats.processing_interval_ms:5000}")
private int statsProcessingInterval;
public DefaultJobManager(JobService jobService, JobStatsService jobStatsService, TbCoreQueueFactory queueFactory, List<JobProcessor> jobProcessors) {
this.jobService = jobService; this.jobService = jobService;
this.queueFactory = queueFactory; this.jobStatsService = jobStatsService;
this.jobProcessors = jobProcessors.stream().collect(Collectors.toMap(JobProcessor::getType, Function.identity())); this.jobProcessors = jobProcessors.stream().collect(Collectors.toMap(JobProcessor::getType, Function.identity()));
this.taskProducers = Arrays.stream(JobType.values()).collect(Collectors.toMap(Function.identity(), queueFactory::createTaskProducer)); this.taskProducers = Arrays.stream(JobType.values()).collect(Collectors.toMap(Function.identity(), queueFactory::createTaskProducer));
this.consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("task-result-consumer")); this.consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("task-result-consumer"));
this.taskResultConsumer = QueueConsumerManager.<TbProtoQueueMsg<TaskResultProto>>builder() // fixme: should be consumer per partition this.taskResultConsumer = QueueConsumerManager.<TbProtoQueueMsg<JobStatsMsg>>builder()
.name("tasks-results") .name("job-stats")
.msgPackProcessor(this::processResults) .msgPackProcessor(this::processStats)
.pollInterval(125) .pollInterval(125)
.consumerCreator(queueFactory::createTaskResultConsumer) .consumerCreator(queueFactory::createJobStatsConsumer)
.consumerExecutor(consumerExecutor) .consumerExecutor(consumerExecutor)
.build(); .build();
} }
@ -82,10 +90,13 @@ public class DefaultJobManager implements JobManager {
} }
@Override @Override
public void submitJob(Job job) { public Job submitJob(Job job) {
job = jobService.createJob(job.getTenantId(), job); job = jobService.createJob(job.getTenantId(), job);
log.info("Submitting job: {}", job); log.info("Submitting job: {}", job);
jobProcessors.get(job.getType()).process(job, this::submitTask);
int tasksCount = jobProcessors.get(job.getType()).process(job, this::submitTask);
jobStatsService.reportAllTasksSubmitted(job.getId(), tasksCount);
return job;
} }
private void submitTask(Task task) { private void submitTask(Task task) {
@ -110,21 +121,34 @@ public class DefaultJobManager implements JobManager {
} }
@SneakyThrows @SneakyThrows
private void processResults(List<TbProtoQueueMsg<TaskResultProto>> msgs, TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> consumer) { private void processStats(List<TbProtoQueueMsg<JobStatsMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<JobStatsMsg>> consumer) {
Map<JobId, List<TaskResult>> results = msgs.stream() Map<JobId, JobStats> stats = new HashMap<>();
.map(msg -> JacksonUtil.fromString(msg.getValue().getValue(), TaskResult.class))
.collect(Collectors.groupingBy(TaskResult::getJobId)); for (TbProtoQueueMsg<JobStatsMsg> msg : msgs) {
results.forEach((jobId, taskResults) -> { JobStatsMsg statsMsg = msg.getValue();
JobId jobId = new JobId(new UUID(statsMsg.getJobIdMSB(), statsMsg.getJobIdLSB()));
JobStats jobStats = stats.computeIfAbsent(jobId, JobStats::new);
if (statsMsg.hasTaskResult()) {
TaskResult taskResult = JacksonUtil.fromString(statsMsg.getTaskResult().getValue(), TaskResult.class);
jobStats.getTaskResults().add(taskResult);
}
if (statsMsg.hasTotalTasksCount()) {
jobStats.setTotalTasksCount(statsMsg.getTotalTasksCount());
}
}
stats.forEach((jobId, jobStats) -> {
try { try {
log.info("[{}] Processing task results: {}", jobId, taskResults); log.info("[{}] Processing job stats: {}", jobId, stats);
jobService.reportTaskResults(jobId, taskResults); jobService.processStats(jobId, jobStats);
} catch (Exception e) { } catch (Exception e) {
log.warn("Failed to report task results for job {}: {}", jobId, taskResults, e); log.warn("Failed to process job stats for {}: {}", jobId, jobStats, e);
} }
}); });
consumer.commit(); consumer.commit();
Thread.sleep(5000); Thread.sleep(statsProcessingInterval);
} }
@PreDestroy @PreDestroy

64
application/src/main/java/org/thingsboard/server/service/job/DummyJobProcessor.java

@ -0,0 +1,64 @@
/**
* 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 lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.job.DummyJobConfiguration;
import org.thingsboard.server.common.data.job.DummyTask;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.Task;
import java.util.List;
import java.util.function.Consumer;
@Component
@RequiredArgsConstructor
public class DummyJobProcessor extends JobProcessor {
@Override
public int process(Job job, Consumer<Task> taskConsumer) {
DummyJobConfiguration configuration = job.getConfiguration();
for (int number = 1; number <= configuration.getSuccessfulTasksCount(); number++) {
taskConsumer.accept(createTask(job, configuration, number, null));
}
if (configuration.getErrors() != null) {
for (int number = 1; number <= configuration.getFailedTasksCount(); number++) {
taskConsumer.accept(createTask(job, configuration, number, configuration.getErrors()));
}
}
return configuration.getSuccessfulTasksCount() + configuration.getFailedTasksCount();
}
private Task createTask(Job job, DummyJobConfiguration configuration, int number, List<String> errors) {
return DummyTask.builder()
.tenantId(job.getTenantId())
.jobId(job.getId())
.key("Task " + number)
.retries(configuration.getRetries())
.number(number)
.processingTimeMs(configuration.getTaskProcessingTimeMs())
.errors(errors)
.build();
}
@Override
public JobType getType() {
return JobType.DUMMY;
}
}

2
application/src/main/java/org/thingsboard/server/service/job/JobManager.java

@ -19,6 +19,6 @@ import org.thingsboard.server.common.data.job.Job;
public interface JobManager { public interface JobManager {
void submitJob(Job job); Job submitJob(Job job);
} }

2
application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java

@ -23,7 +23,7 @@ import java.util.function.Consumer;
public abstract class JobProcessor { public abstract class JobProcessor {
public abstract void process(Job job, Consumer<Task> taskConsumer); public abstract int process(Job job, Consumer<Task> taskConsumer);
public abstract JobType getType(); public abstract JobType getType();

44
application/src/main/java/org/thingsboard/server/service/job/task/DummyTaskProcessor.java

@ -0,0 +1,44 @@
/**
* 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.task;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.job.DummyTask;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.queue.task.TaskProcessor;
@Component
@RequiredArgsConstructor
public class DummyTaskProcessor extends TaskProcessor<DummyTask> {
@Override
protected void process(DummyTask task) throws Exception {
if (task.getProcessingTimeMs() > 0) {
Thread.sleep(task.getProcessingTimeMs());
}
if (task.getErrors() != null && task.getAttempt() <= task.getErrors().size()) {
String error = task.getErrors().get(task.getAttempt() - 1);
throw new RuntimeException(error);
}
}
@Override
public JobType getJobType() {
return JobType.DUMMY;
}
}

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

@ -1883,6 +1883,10 @@ queue:
enabled: "${TB_QUEUE_EDGE_STATS_ENABLED:true}" enabled: "${TB_QUEUE_EDGE_STATS_ENABLED:true}"
# Statistics printing interval for Edge services # Statistics printing interval for Edge services
print-interval-ms: "${TB_QUEUE_EDGE_STATS_PRINT_INTERVAL_MS:60000}" print-interval-ms: "${TB_QUEUE_EDGE_STATS_PRINT_INTERVAL_MS:60000}"
tasks:
stats:
# Interval in milliseconds to process job stats
processing_interval_ms: "${TB_QUEUE_TASKS_STATS_PROCESSING_INTERVAL_MS:5000}"
# Event configuration parameters # Event configuration parameters
event: event:

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

@ -0,0 +1,125 @@
/**
* 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 com.fasterxml.jackson.core.type.TypeReference;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.job.DummyJobConfiguration;
import org.thingsboard.server.common.data.job.Job;
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.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
@DaoSqlTest
@TestPropertySource(properties = {
"queue.tasks.stats.processing_interval_ms=0"
})
public class JobManagerTest extends AbstractControllerTest {
@Autowired
private JobManager jobManager;
@Before
public void setUp() throws Exception {
loginTenantAdmin();
}
@After
public void tearDown() throws Exception {
}
@Test
public void testSubmitJob_allTasksSuccessful() {
int tasksCount = 5;
JobId jobId = jobManager.submitJob(Job.builder()
.tenantId(tenantId)
.type(JobType.DUMMY)
.key("test-job")
.description("test job")
.configuration(DummyJobConfiguration.builder()
.successfulTasksCount(tasksCount)
.taskProcessingTimeMs(1000)
.build())
.build()).getId();
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {
Job job = findJobById(jobId);
assertThat(job.getStatus()).isEqualTo(JobStatus.RUNNING);
assertThat(job.getResult().getSuccessfulCount()).isBetween(1, tasksCount - 1);
assertThat(job.getResult().getTotalCount()).isEqualTo(tasksCount);
});
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {
Job job = findJobById(jobId);
assertThat(job.getStatus()).isEqualTo(JobStatus.COMPLETED);
assertThat(job.getResult().getSuccessfulCount()).isEqualTo(tasksCount);
assertThat(job.getResult().getFailures()).isEmpty();
});
}
@Test
public void testSubmitJob_someTasksPermanentlyFailed() {
int successfulTasks = 3;
int failedTasks = 2;
JobId jobId = jobManager.submitJob(Job.builder()
.tenantId(tenantId)
.type(JobType.DUMMY)
.key("test-job")
.description("test job")
.configuration(DummyJobConfiguration.builder()
.successfulTasksCount(successfulTasks)
.failedTasksCount(failedTasks)
.errors(List.of("error1", "error2", "error3"))
.retries(2)
.taskProcessingTimeMs(100)
.build())
.build()).getId();
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> {
Job job = findJobById(jobId);
assertThat(job.getStatus()).isEqualTo(JobStatus.FAILED);
assertThat(job.getResult().getSuccessfulCount()).isEqualTo(successfulTasks);
assertThat(job.getResult().getFailedCount()).isEqualTo(failedTasks);
assertThat(job.getResult().getTotalCount()).isEqualTo(successfulTasks + failedTasks);
assertThat(job.getResult().getFailures().get("Task 1")).isEqualTo("error3"); // last error
assertThat(job.getResult().getFailures().get("Task 2")).isEqualTo("error3"); // last error
});
}
private Job findJobById(JobId jobId) throws Exception {
return doGet("/api/job/" + jobId, Job.class);
}
private List<Job> findJobs() throws Exception {
return doGetTypedWithPageLink("/api/jobs?", new TypeReference<PageData<Job>>() {}, new PageLink(100, 0)).getData();
}
}

10
common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java

@ -17,18 +17,18 @@ package org.thingsboard.server.dao.task;
import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobStats;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.job.Job;
import java.util.List;
public interface JobService { public interface JobService {
Job createJob(TenantId tenantId, Job job); Job createJob(TenantId tenantId, Job job);
void reportTaskResults(JobId jobId, List<TaskResult> results); Job findJobById(TenantId tenantId, JobId jobId);
void processStats(JobId jobId, JobStats jobStats);
PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink); PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink);

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

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.common.data.job; package org.thingsboard.server.common.data.job;
import jakarta.validation.constraints.NotNull;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
import lombok.Builder; import lombok.Builder;
import lombok.Data; import lombok.Data;
@ -27,6 +28,7 @@ import org.thingsboard.server.common.data.cf.CalculatedField;
@Builder @Builder
public class CfReprocessingJobConfiguration implements JobConfiguration { public class CfReprocessingJobConfiguration implements JobConfiguration {
@NotNull
private CalculatedField calculatedField; private CalculatedField calculatedField;
private long startTs; private long startTs;
private long endTs; private long endTs;

16
common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java

@ -15,20 +15,17 @@
*/ */
package org.thingsboard.server.common.data.job; package org.thingsboard.server.common.data.job;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode; import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import lombok.experimental.SuperBuilder;
import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
@Data @Data
@AllArgsConstructor
@NoArgsConstructor @NoArgsConstructor
@EqualsAndHashCode(callSuper = true) @EqualsAndHashCode(callSuper = true)
@SuperBuilder
public class CfReprocessingTask extends Task { public class CfReprocessingTask extends Task {
private CalculatedField calculatedField; private CalculatedField calculatedField;
@ -36,15 +33,6 @@ public class CfReprocessingTask extends Task {
private long startTs; private long startTs;
private long endTs; private long endTs;
@Builder
public CfReprocessingTask(TenantId tenantId, JobId jobId, String key, CalculatedField calculatedField, EntityId entityId, long startTs, long endTs) {
super(tenantId, jobId, key);
this.calculatedField = calculatedField;
this.entityId = entityId;
this.startTs = startTs;
this.endTs = endTs;
}
@Override @Override
public JobType getJobType() { public JobType getJobType() {
return JobType.CF_REPROCESSING; return JobType.CF_REPROCESSING;

42
common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobConfiguration.java

@ -0,0 +1,42 @@
/**
* 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.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.List;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class DummyJobConfiguration implements JobConfiguration {
private long taskProcessingTimeMs;
private int successfulTasksCount;
private int failedTasksCount;
private List<String> errors;
private int retries;
@Override
public JobType getType() {
return JobType.DUMMY;
}
}

25
common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobResult.java

@ -0,0 +1,25 @@
/**
* 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;
public class DummyJobResult extends JobResult {
@Override
public JobType getJobType() {
return JobType.DUMMY;
}
}

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

@ -0,0 +1,40 @@
/**
* 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.Data;
import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import lombok.experimental.SuperBuilder;
import java.util.List;
@Data
@NoArgsConstructor
@EqualsAndHashCode(callSuper = true)
@SuperBuilder
public class DummyTask extends Task {
private int number;
private long processingTimeMs;
private List<String> errors; // errors for each attempt
@Override
public JobType getJobType() {
return JobType.DUMMY;
}
}

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

@ -15,6 +15,8 @@
*/ */
package org.thingsboard.server.common.data.job; package org.thingsboard.server.common.data.job;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import lombok.Builder; import lombok.Builder;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode; import lombok.EqualsAndHashCode;
@ -29,22 +31,30 @@ import org.thingsboard.server.common.data.id.TenantId;
@EqualsAndHashCode(callSuper = true) @EqualsAndHashCode(callSuper = true)
public class Job extends BaseData<JobId> implements HasTenantId { public class Job extends BaseData<JobId> implements HasTenantId {
@NotNull
private TenantId tenantId; private TenantId tenantId;
@NotNull
private JobType type; private JobType type;
@NotBlank
private String key; private String key;
@NotBlank
private String description;
private JobStatus status; private JobStatus status;
@NotNull
private JobConfiguration configuration; private JobConfiguration configuration;
private JobResult result; private JobResult result;
@Builder @Builder
public Job(TenantId tenantId, JobType type, String key, JobConfiguration configuration) { public Job(TenantId tenantId, JobType type, String key, String description, JobConfiguration configuration) {
this.tenantId = tenantId; this.tenantId = tenantId;
this.type = type; this.type = type;
this.key = key; this.key = key;
this.description = description;
this.configuration = configuration; this.configuration = configuration;
this.status = JobStatus.PENDING; this.status = JobStatus.PENDING;
this.result = switch (type) { this.result = switch (type) {
case CF_REPROCESSING -> new CfReprocessingJobResult(); case CF_REPROCESSING -> new CfReprocessingJobResult();
case DUMMY -> new DummyJobResult();
}; };
} }

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

@ -24,6 +24,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo;
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type") @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type")
@JsonSubTypes({ @JsonSubTypes({
@Type(name = "CF_REPROCESSING", value = CfReprocessingJobConfiguration.class), @Type(name = "CF_REPROCESSING", value = CfReprocessingJobConfiguration.class),
@Type(name = "DUMMY", value = DummyJobConfiguration.class),
}) })
public interface JobConfiguration { public interface JobConfiguration {

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

@ -17,8 +17,10 @@ package org.thingsboard.server.common.data.job;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.HashMap; import java.util.HashMap;
import java.util.Map; import java.util.Map;
@ -26,14 +28,16 @@ import java.util.Map;
@JsonIgnoreProperties(ignoreUnknown = true) @JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType")
@JsonSubTypes({ @JsonSubTypes({
@JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class), @Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class),
@Type(name = "DUMMY", value = DummyJobResult.class)
}) })
@Data @Data
@NoArgsConstructor
public abstract class JobResult { public abstract class JobResult {
private int successfulCount; private int successfulCount;
private int failedCount; private int failedCount;
private int totalCount; private Integer totalCount = null; // set when all tasks are submitted
private Map<String, String> failures = new HashMap<>(); private Map<String, String> failures = new HashMap<>();
public abstract JobType getJobType(); public abstract JobType getJobType();

29
common/data/src/main/java/org/thingsboard/server/common/data/job/JobStats.java

@ -0,0 +1,29 @@
/**
* 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.Data;
import org.thingsboard.server.common.data.id.JobId;
import java.util.ArrayList;
import java.util.List;
@Data
public class JobStats {
private final JobId jobId;
private final List<TaskResult> taskResults = new ArrayList<>();
private Integer totalTasksCount;
}

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

@ -17,7 +17,8 @@ package org.thingsboard.server.common.data.job;
public enum JobType { public enum JobType {
CF_REPROCESSING; CF_REPROCESSING,
DUMMY;
public String getTasksTopic() { public String getTasksTopic() {
return "tasks." + name().toLowerCase(); return "tasks." + name().toLowerCase();

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

@ -17,8 +17,11 @@ package org.thingsboard.server.common.data.job;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.AllArgsConstructor;
import lombok.Data; import lombok.Data;
import lombok.experimental.SuperBuilder;
import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -26,19 +29,17 @@ import org.thingsboard.server.common.data.id.TenantId;
@JsonIgnoreProperties(ignoreUnknown = true) @JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType")
@JsonSubTypes({ @JsonSubTypes({
@JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class), @Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class),
@Type(name = "DUMMY", value = DummyTask.class)
}) })
@SuperBuilder
@AllArgsConstructor
public abstract class Task { public abstract class Task {
private TenantId tenantId; private TenantId tenantId;
private JobId jobId; private JobId jobId;
private String key; private String key;
private int retries;
public Task(TenantId tenantId, JobId jobId, String key) {
this.tenantId = tenantId;
this.jobId = jobId;
this.key = key;
}
public Task() { public Task() {
} }

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

@ -19,8 +19,6 @@ import lombok.AllArgsConstructor;
import lombok.Builder; import lombok.Builder;
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
@Data @Data
@AllArgsConstructor @AllArgsConstructor
@ -28,8 +26,6 @@ import org.thingsboard.server.common.data.id.TenantId;
@Builder @Builder
public class TaskResult { public class TaskResult {
private TenantId tenantId;
private JobId jobId;
private boolean success; private boolean success;
private TaskFailure failure; private TaskFailure failure;

7
common/proto/src/main/proto/queue.proto

@ -1851,6 +1851,13 @@ message TaskProto {
string value = 1; // fixme: TMP, make more efficient string value = 1; // fixme: TMP, make more efficient
} }
message JobStatsMsg {
int64 jobIdMSB = 1;
int64 jobIdLSB = 2;
optional TaskResultProto taskResult = 3;
optional int32 totalTasksCount = 4;
}
message TaskResultProto { message TaskResultProto {
string value = 1; // fixme: TMP, make more efficient string value = 1; // fixme: TMP, make more efficient
} }

9
common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java

@ -28,6 +28,7 @@ import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueAdmin;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
@ -270,13 +271,13 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
} }
@Override @Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultProducer() { public TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsProducer() {
return new InMemoryTbQueueProducer<>(storage, "tasks.results"); return new InMemoryTbQueueProducer<>(storage, "jobs.stats");
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultConsumer() { public TbQueueConsumer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsConsumer() {
return new InMemoryTbQueueConsumer<>(storage, "tasks.results"); return new InMemoryTbQueueConsumer<>(storage, "jobs.stats");
} }
@Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}") @Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}")

22
common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java

@ -30,8 +30,8 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -670,23 +670,23 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
} }
@Override @Override
public TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> createTaskResultProducer() { public TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsProducer() {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TaskResultProto>>builder() return TbKafkaProducerTemplate.<TbProtoQueueMsg<JobStatsMsg>>builder()
.clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) .clientId("job-stats-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName("tasks.results")) .defaultTopic(topicService.buildTopicName("jobs.stats"))
.settings(kafkaSettings) .settings(kafkaSettings)
.admin(tasksAdmin) .admin(tasksAdmin)
.build(); .build();
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> createTaskResultConsumer() { public TbQueueConsumer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsConsumer() {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TaskResultProto>>builder() return TbKafkaConsumerTemplate.<TbProtoQueueMsg<JobStatsMsg>>builder()
.settings(kafkaSettings) .settings(kafkaSettings)
.topic(topicService.buildTopicName("tasks.results")) .topic(topicService.buildTopicName("jobs.stats"))
.clientId("task-result-consumer-" + serviceInfoProvider.getServiceId()) .clientId("job-stats-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName("task-result-consumer-group")) .groupId(topicService.buildTopicName("job-stats-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskResultProto.parseFrom(msg.getData()), msg.getHeaders())) .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), JobStatsMsg.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin) .admin(tasksAdmin)
.statsService(consumerStatsService) .statsService(consumerStatsService)
.build(); .build();

37
common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java

@ -25,8 +25,8 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
@ -536,13 +536,36 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultConsumer() { public TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType) {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TransportProtos.TaskResultProto>>builder() return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TaskProto>>builder()
.settings(kafkaSettings) .settings(kafkaSettings)
.topic(topicService.buildTopicName("tasks.results")) .topic(topicService.buildTopicName(jobType.getTasksTopic()))
.clientId("task-result-consumer-" + serviceInfoProvider.getServiceId()) .clientId(jobType.name().toLowerCase() + "-task-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName("task-result-consumer-group")) .groupId(topicService.buildTopicName(jobType.name().toLowerCase() + "-task-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.TaskResultProto.parseFrom(msg.getData()), msg.getHeaders())) .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskProto.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin)
.statsService(consumerStatsService)
.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsProducer() {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<JobStatsMsg>>builder()
.clientId("job-stats-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName("jobs.stats"))
.settings(kafkaSettings)
.admin(tasksAdmin)
.build();
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsConsumer() {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<JobStatsMsg>>builder()
.settings(kafkaSettings)
.topic(topicService.buildTopicName("jobs.stats"))
.clientId("job-stats-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName("job-stats-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), JobStatsMsg.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin) .admin(tasksAdmin)
.statsService(consumerStatsService) .statsService(consumerStatsService)
.build(); .build();

10
common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java

@ -27,9 +27,9 @@ import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
@ -433,10 +433,10 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
} }
@Override @Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultProducer() { public TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsProducer() {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TransportProtos.TaskResultProto>>builder() return TbKafkaProducerTemplate.<TbProtoQueueMsg<JobStatsMsg>>builder()
.clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) .clientId("job-stats-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName("tasks.results")) .defaultTopic(topicService.buildTopicName("jobs.stats"))
.settings(kafkaSettings) .settings(kafkaSettings)
.admin(tasksAdmin) .admin(tasksAdmin)
.build(); .build();

4
common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java

@ -16,8 +16,8 @@
package org.thingsboard.server.queue.provider; package org.thingsboard.server.queue.provider;
import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
@ -26,6 +26,6 @@ public interface TaskProcessorQueueFactory {
TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType); TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType);
TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> createTaskResultProducer(); TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsProducer();
} }

6
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java

@ -19,8 +19,8 @@ import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -47,7 +47,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
* Responsible for initialization of various Producers and Consumers used by TB Core Node. * Responsible for initialization of various Producers and Consumers used by TB Core Node.
* Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable * Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable
*/ */
public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory { public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory, TaskProcessorQueueFactory {
/** /**
* Used to push messages to instances of TB Transport Service * Used to push messages to instances of TB Transport Service
@ -170,6 +170,6 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous
TbQueueProducer<TbProtoQueueMsg<TaskProto>> createTaskProducer(JobType jobType); TbQueueProducer<TbProtoQueueMsg<TaskProto>> createTaskProducer(JobType jobType);
TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> createTaskResultConsumer(); TbQueueConsumer<TbProtoQueueMsg<JobStatsMsg>> createJobStatsConsumer();
} }

8
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java

@ -18,6 +18,7 @@ package org.thingsboard.server.queue.provider;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -53,6 +54,7 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider {
private TbQueueProducer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> toHousekeeper; private TbQueueProducer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> toHousekeeper;
private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> toCalculatedFields; private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> toCalculatedFields;
private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> toCalculatedFieldNotifications; private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> toCalculatedFieldNotifications;
private TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> jobStatsProducer;
public TbCoreQueueProducerProvider(TbCoreQueueFactory tbQueueProvider) { public TbCoreQueueProducerProvider(TbCoreQueueFactory tbQueueProvider) {
this.tbQueueProvider = tbQueueProvider; this.tbQueueProvider = tbQueueProvider;
@ -73,6 +75,7 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider {
this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer(); this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer();
this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer(); this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer();
this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer(); this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer();
this.jobStatsProducer = tbQueueProvider.createJobStatsProducer();
} }
@Override @Override
@ -140,4 +143,9 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider {
return toCalculatedFieldNotifications; return toCalculatedFieldNotifications;
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> getJobStatsProducer() {
return jobStatsProducer;
}
} }

3
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.queue.provider; package org.thingsboard.server.queue.provider;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -97,4 +98,6 @@ public interface TbQueueProducerProvider {
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> getCalculatedFieldsNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> getCalculatedFieldsNotificationsMsgProducer();
TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> getJobStatsProducer();
} }

8
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java

@ -18,6 +18,7 @@ package org.thingsboard.server.queue.provider;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -51,6 +52,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider {
private TbQueueProducer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> toEdgeEvents; private TbQueueProducer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> toEdgeEvents;
private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> toCalculatedFields; private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> toCalculatedFields;
private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> toCalculatedFieldNotifications; private TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> toCalculatedFieldNotifications;
private TbQueueProducer<TbProtoQueueMsg<TransportProtos.JobStatsMsg>> jobStatsProducer;
public TbRuleEngineProducerProvider(TbRuleEngineQueueFactory tbQueueProvider) { public TbRuleEngineProducerProvider(TbRuleEngineQueueFactory tbQueueProvider) {
this.tbQueueProvider = tbQueueProvider; this.tbQueueProvider = tbQueueProvider;
@ -70,6 +72,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider {
this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer(); this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer();
this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer(); this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer();
this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer(); this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer();
this.jobStatsProducer = tbQueueProvider.createJobStatsProducer();
} }
@Override @Override
@ -137,4 +140,9 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider {
return toCalculatedFieldNotifications; return toCalculatedFieldNotifications;
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.JobStatsMsg>> getJobStatsProducer() {
return jobStatsProducer;
}
} }

6
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java

@ -121,4 +121,10 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToCalculatedFieldNotificationMsg>> getCalculatedFieldsNotificationsMsgProducer() { public TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToCalculatedFieldNotificationMsg>> getCalculatedFieldsNotificationsMsgProducer() {
throw new RuntimeException("Not Implemented! Should not be used by Transport!"); throw new RuntimeException("Not Implemented! Should not be used by Transport!");
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.JobStatsMsg>> getJobStatsProducer() {
throw new RuntimeException("Not Implemented! Should not be used by Transport!");
}
} }

5
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java

@ -118,4 +118,9 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider
throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!");
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.JobStatsMsg>> getJobStatsProducer() {
throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!");
}
} }

63
common/queue/src/main/java/org/thingsboard/server/queue/task/JobStatsService.java

@ -0,0 +1,63 @@
/**
* 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.queue.task;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
@Lazy
@Service
@Slf4j
@RequiredArgsConstructor
public class JobStatsService {
private final TbQueueProducerProvider producerProvider;
public void reportTaskResult(JobId jobId, TaskResult result) {
report(jobId, JobStatsMsg.newBuilder()
.setTaskResult(TaskResultProto.newBuilder()
.setValue(JacksonUtil.toString(result))
.build()));
}
public void reportAllTasksSubmitted(JobId jobId, int tasksCount) {
report(jobId, JobStatsMsg.newBuilder()
.setTotalTasksCount(tasksCount));
}
private void report(JobId jobId, JobStatsMsg.Builder statsMsg) {
log.info("[{}] Reporting: {}", jobId, statsMsg);
statsMsg.setJobIdMSB(jobId.getId().getMostSignificantBits())
.setJobIdLSB(jobId.getId().getLeastSignificantBits());
TbProtoQueueMsg<JobStatsMsg> msg = new TbProtoQueueMsg<>(jobId.getId(), statsMsg.build());
TbQueueProducer<TbProtoQueueMsg<JobStatsMsg>> producer = producerProvider.getJobStatsProducer();
producer.send(TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(), msg, TbQueueCallback.EMPTY);
}
}

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

@ -25,12 +25,8 @@ import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.Task; import org.thingsboard.server.common.data.job.Task;
import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.job.TaskResult.TaskFailure; import org.thingsboard.server.common.data.job.TaskResult.TaskFailure;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto;
import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager;
import org.thingsboard.server.queue.provider.TaskProcessorQueueFactory; import org.thingsboard.server.queue.provider.TaskProcessorQueueFactory;
@ -45,22 +41,22 @@ public abstract class TaskProcessor<T extends Task> {
@Autowired @Autowired
private TaskProcessorQueueFactory queueFactory; private TaskProcessorQueueFactory queueFactory;
@Autowired
private JobStatsService statsService;
private QueueConsumerManager<TbProtoQueueMsg<TaskProto>> taskConsumer; private QueueConsumerManager<TbProtoQueueMsg<TaskProto>> taskConsumer;
private TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> taskResultProducer;
private ExecutorService consumerExecutor; private ExecutorService consumerExecutor;
@PostConstruct @PostConstruct
public void init() { public void init() {
consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-consumer")); consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-consumer"));
taskConsumer = QueueConsumerManager.<TbProtoQueueMsg<TaskProto>>builder() // fixme: should be consumer per partition taskConsumer = QueueConsumerManager.<TbProtoQueueMsg<TaskProto>>builder() // fixme: should be consumer per partition
.name(getJobType().name() + "-tasks") .name(getJobType().name().toLowerCase() + "-tasks")
.msgPackProcessor(this::processMsgs) .msgPackProcessor(this::processMsgs)
.pollInterval(125) .pollInterval(125)
.consumerCreator(() -> queueFactory.createTaskConsumer(getJobType())) .consumerCreator(() -> queueFactory.createTaskConsumer(getJobType()))
.consumerExecutor(consumerExecutor) .consumerExecutor(consumerExecutor)
.build(); .build();
taskResultProducer = queueFactory.createTaskResultProducer();
} }
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
@ -92,7 +88,7 @@ public abstract class TaskProcessor<T extends Task> {
reportSuccess(task); reportSuccess(task);
} catch (Exception e) { } catch (Exception e) {
log.error("Failed to process task (attempt {}): {}", task.getAttempt(), task, e); log.error("Failed to process task (attempt {}): {}", task.getAttempt(), task, e);
if (task.getAttempt() < 3) { if (task.getAttempt() <= task.getRetries()) {
processTask(task); processTask(task);
} else { } else {
reportFailure(task, e); reportFailure(task, e);
@ -102,34 +98,19 @@ public abstract class TaskProcessor<T extends Task> {
private void reportSuccess(Task task) { private void reportSuccess(Task task) {
TaskResult result = TaskResult.builder() TaskResult result = TaskResult.builder()
.tenantId(task.getTenantId())
.jobId(task.getJobId())
.success(true) .success(true)
.build(); .build();
reportResult(result); statsService.reportTaskResult(task.getJobId(), result);
} }
private void reportFailure(Task task, Throwable error) { private void reportFailure(Task task, Throwable error) {
TaskResult result = TaskResult.builder() TaskResult result = TaskResult.builder()
.tenantId(task.getTenantId())
.jobId(task.getJobId())
.failure(TaskFailure.builder() .failure(TaskFailure.builder()
.error(error.getMessage()) .error(error.getMessage())
.task(task) .task(task)
.build()) .build())
.build(); .build();
reportResult(result); statsService.reportTaskResult(task.getJobId(), result);
}
private void reportResult(TaskResult result) {
log.info("Reporting result: {}", result);
TaskResultProto resultProto = TaskResultProto.newBuilder()
.setValue(JacksonUtil.toString(result))
.build();
TbProtoQueueMsg<TaskResultProto> msg = new TbProtoQueueMsg<>(result.getJobId().getId(), resultProto);
taskResultProducer.send(TopicPartitionInfo.builder()
.topic(taskResultProducer.getDefaultTopic())
.build(), msg, TbQueueCallback.EMPTY);
} }
protected abstract void process(T task) throws Exception; protected abstract void process(T task) throws Exception;

1
dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java

@ -745,6 +745,7 @@ public class ModelConstants {
public static final String JOB_TABLE_NAME = "job"; public static final String JOB_TABLE_NAME = "job";
public static final String JOB_TYPE_PROPERTY = "type"; public static final String JOB_TYPE_PROPERTY = "type";
public static final String JOB_KEY_PROPERTY = "key"; public static final String JOB_KEY_PROPERTY = "key";
public static final String JOB_DESCRIPTION_PROPERTY = "description";
public static final String JOB_STATUS_PROPERTY = "status"; public static final String JOB_STATUS_PROPERTY = "status";
public static final String JOB_CONFIGURATION_PROPERTY = "configuration"; public static final String JOB_CONFIGURATION_PROPERTY = "configuration";
public static final String JOB_RESULT_PROPERTY = "result"; public static final String JOB_RESULT_PROPERTY = "result";

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java

@ -56,6 +56,9 @@ public class JobEntity extends BaseSqlEntity<Job> {
@Column(name = ModelConstants.JOB_KEY_PROPERTY, nullable = false) @Column(name = ModelConstants.JOB_KEY_PROPERTY, nullable = false)
private String key; private String key;
@Column(name = ModelConstants.JOB_DESCRIPTION_PROPERTY, nullable = false)
private String description;
@Enumerated(EnumType.STRING) @Enumerated(EnumType.STRING)
@Column(name = ModelConstants.JOB_STATUS_PROPERTY, nullable = false) @Column(name = ModelConstants.JOB_STATUS_PROPERTY, nullable = false)
private JobStatus status; private JobStatus status;
@ -74,6 +77,7 @@ public class JobEntity extends BaseSqlEntity<Job> {
this.tenantId = getTenantUuid(job.getTenantId()); this.tenantId = getTenantUuid(job.getTenantId());
this.type = job.getType(); this.type = job.getType();
this.key = job.getKey(); this.key = job.getKey();
this.description = job.getDescription();
this.status = job.getStatus(); this.status = job.getStatus();
this.configuration = toJson(job.getConfiguration()); this.configuration = toJson(job.getConfiguration());
this.result = toJson(job.getResult()); this.result = toJson(job.getResult());
@ -87,6 +91,7 @@ public class JobEntity extends BaseSqlEntity<Job> {
job.setTenantId(getTenantId(tenantId)); job.setTenantId(getTenantId(tenantId));
job.setType(type); job.setType(type);
job.setKey(key); job.setKey(key);
job.setDescription(description);
job.setStatus(status); job.setStatus(status);
job.setConfiguration(fromJson(configuration, JobConfiguration.class)); job.setConfiguration(fromJson(configuration, JobConfiguration.class));
job.setResult(fromJson(result, JobResult.class)); job.setResult(fromJson(result, JobResult.class));

21
dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java

@ -23,14 +23,22 @@ import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query; import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param; import org.springframework.data.repository.query.Param;
import org.springframework.stereotype.Repository; import org.springframework.stereotype.Repository;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.dao.model.sql.JobEntity; import org.thingsboard.server.dao.model.sql.JobEntity;
import java.util.List;
import java.util.UUID; import java.util.UUID;
@Repository @Repository
public interface JobRepository extends JpaRepository<JobEntity, UUID> { public interface JobRepository extends JpaRepository<JobEntity, UUID> {
Page<JobEntity> findByTenantId(UUID tenantId, Pageable pageable); @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);
@Modifying @Modifying
@Transactional @Transactional
@ -45,7 +53,8 @@ public interface JobRepository extends JpaRepository<JobEntity, UUID> {
RETURNING ((result->>'successfulCount')::int + :count) RETURNING ((result->>'successfulCount')::int + :count)
+ (result->>'failedCount')::int = (result->>'totalCount')::int + (result->>'failedCount')::int = (result->>'totalCount')::int
""", nativeQuery = true) """, nativeQuery = true)
boolean reportTaskSuccess(@Param("jobId") UUID jobId, @Param("count") int count); boolean reportTaskSuccess(@Param("jobId") UUID jobId,
@Param("count") int count);
@Modifying @Modifying
@Transactional @Transactional
@ -64,6 +73,12 @@ public interface JobRepository extends JpaRepository<JobEntity, UUID> {
RETURNING ((result->>'failedCount')::int + 1) + (result->>'successfulCount')::int RETURNING ((result->>'failedCount')::int + 1) + (result->>'successfulCount')::int
= (result->>'totalCount')::int = (result->>'totalCount')::int
""", nativeQuery = true) """, nativeQuery = true)
boolean reportTaskFailure(@Param("jobId") UUID jobId, @Param("taskKey") String taskKey, @Param("error") String error); boolean reportTaskFailure(@Param("jobId") UUID jobId,
@Param("taskKey") String taskKey,
@Param("error") String error);
boolean existsByKeyAndStatusIn(String key, List<JobStatus> statuses);
boolean existsByTenantIdAndTypeAndStatusIn(UUID tenantId, JobType type, List<JobStatus> statuses);
} }

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

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.sql.task; package org.thingsboard.server.dao.sql.task;
import com.google.common.base.Strings;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@ -22,6 +23,8 @@ import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job; import org.thingsboard.server.common.data.job.Job;
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.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
@ -30,6 +33,7 @@ import org.thingsboard.server.dao.sql.JpaAbstractDao;
import org.thingsboard.server.dao.task.JobDao; import org.thingsboard.server.dao.task.JobDao;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.Arrays;
import java.util.UUID; import java.util.UUID;
@Component @Component
@ -41,7 +45,7 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
@Override @Override
public PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink) { public PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink) {
return DaoUtil.toPageData(jobRepository.findByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink))); return DaoUtil.toPageData(jobRepository.findByTenantIdAndSearchText(tenantId.getId(), Strings.emptyToNull(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink)));
} }
@Override @Override
@ -54,6 +58,16 @@ public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao
return jobRepository.reportTaskFailure(jobId.getId(), taskKey, error); return jobRepository.reportTaskFailure(jobId.getId(), taskKey, error);
} }
@Override
public boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses) {
return jobRepository.existsByKeyAndStatusIn(key, Arrays.stream(statuses).toList());
}
@Override
public boolean existsByTenantIdAndTypeAndStatusOneOf(TenantId tenantId, JobType type, JobStatus... statuses) {
return jobRepository.existsByTenantIdAndTypeAndStatusIn(tenantId.getId(), type, Arrays.stream(statuses).toList());
}
@Override @Override
public EntityType getEntityType() { public EntityType getEntityType() {
return EntityType.JOB; return EntityType.JOB;

40
dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java

@ -22,13 +22,14 @@ import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job; import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobResult; 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.JobStatus;
import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.job.TaskResult.TaskFailure; import org.thingsboard.server.common.data.job.TaskResult.TaskFailure;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.exception.DataValidationException;
import java.util.List; import org.thingsboard.server.dao.service.DataValidator;
@Service @Service
@RequiredArgsConstructor @RequiredArgsConstructor
@ -36,14 +37,21 @@ import java.util.List;
public class DefaultJobService implements JobService { public class DefaultJobService implements JobService {
private final JobDao jobDao; private final JobDao jobDao;
private final JobValidator validator = new JobValidator();
@Override @Override
public Job createJob(TenantId tenantId, Job job) { public Job createJob(TenantId tenantId, Job job) {
validator.validate(job, Job::getTenantId);
return jobDao.save(tenantId, job); return jobDao.save(tenantId, job);
} }
@Override @Override
public void reportTaskResults(JobId jobId, List<TaskResult> results) { public Job findJobById(TenantId tenantId, JobId jobId) {
return jobDao.findById(tenantId, jobId.getId());
}
@Override
public void processStats(JobId jobId, JobStats jobStats) {
Job job = jobDao.findById(TenantId.SYS_TENANT_ID, jobId.getId()); Job job = jobDao.findById(TenantId.SYS_TENANT_ID, jobId.getId());
switch (job.getStatus()) { switch (job.getStatus()) {
case PENDING -> { case PENDING -> {
@ -56,7 +64,11 @@ public class DefaultJobService implements JobService {
} }
JobResult jobResult = job.getResult(); JobResult jobResult = job.getResult();
for (TaskResult taskResult : results) { if (jobStats.getTotalTasksCount() != null) {
jobResult.setTotalCount(jobStats.getTotalTasksCount());
}
for (TaskResult taskResult : jobStats.getTaskResults()) {
if (taskResult.isSuccess()) { if (taskResult.isSuccess()) {
jobResult.setSuccessfulCount(jobResult.getSuccessfulCount() + 1); jobResult.setSuccessfulCount(jobResult.getSuccessfulCount() + 1);
} else { } else {
@ -67,7 +79,7 @@ public class DefaultJobService implements JobService {
} }
} }
if (jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) { if (jobResult.getTotalCount() != null && jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) {
if (jobResult.getFailures().isEmpty()) { if (jobResult.getFailures().isEmpty()) {
job.setStatus(JobStatus.COMPLETED); job.setStatus(JobStatus.COMPLETED);
} else { } else {
@ -83,4 +95,22 @@ public class DefaultJobService implements JobService {
return jobDao.findByTenantId(tenantId, pageLink); return jobDao.findByTenantId(tenantId, pageLink);
} }
// todo: cancellation, reprocessing
public class JobValidator extends DataValidator<Job> {
@Override
protected void validateCreate(TenantId tenantId, Job job) {
if (jobDao.existsByTenantIdAndTypeAndStatusOneOf(tenantId, job.getType(), JobStatus.PENDING, JobStatus.RUNNING)) {
throw new DataValidationException("Job of this type is already running");
}
}
@Override
protected Job validateUpdate(TenantId tenantId, Job job) {
throw new IllegalArgumentException("Job can't be updated externally");
}
}
} }

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

@ -18,6 +18,8 @@ package org.thingsboard.server.dao.task;
import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.Job; import org.thingsboard.server.common.data.job.Job;
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.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.Dao;
@ -30,4 +32,8 @@ public interface JobDao extends Dao<Job> {
boolean reportTaskFailure(JobId jobId, String taskKey, String error); boolean reportTaskFailure(JobId jobId, String taskKey, String error);
boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses);
boolean existsByTenantIdAndTypeAndStatusOneOf(TenantId tenantId, JobType type, JobStatus... statuses);
} }

1
dao/src/main/resources/sql/schema-entities.sql

@ -955,6 +955,7 @@ CREATE TABLE IF NOT EXISTS job (
tenant_id uuid NOT NULL, tenant_id uuid NOT NULL,
type varchar NOT NULL, type varchar NOT NULL,
key varchar NOT NULL, key varchar NOT NULL,
description varchar NOT NULL,
status varchar NOT NULL, status varchar NOT NULL,
configuration varchar(1000) NOT NULL, configuration varchar(1000) NOT NULL,
result jsonb result jsonb

Loading…
Cancel
Save