Browse Source

Jobs

pull/13286/head
ViacheslavKlimov 1 year ago
parent
commit
215c5dbb96
  1. 83
      application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java
  2. 136
      application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java
  3. 24
      application/src/main/java/org/thingsboard/server/service/job/JobManager.java
  4. 30
      application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java
  5. 57
      application/src/main/java/org/thingsboard/server/service/job/task/CfReprocessingTaskProcessor.java
  6. 2
      application/src/main/resources/thingsboard.yml
  7. 2
      application/src/test/resources/logback-test.xml
  8. 35
      common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java
  9. 3
      common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java
  10. 2
      common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java
  11. 38
      common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java
  12. 39
      common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java
  13. 25
      common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobResult.java
  14. 53
      common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java
  15. 56
      common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java
  16. 32
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java
  17. 41
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java
  18. 24
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobStatus.java
  19. 26
      common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java
  20. 50
      common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java
  21. 45
      common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java
  22. 8
      common/proto/src/main/proto/queue.proto
  23. 5
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java
  24. 22
      common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java
  25. 51
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java
  26. 28
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java
  27. 29
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java
  28. 31
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java
  29. 7
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java
  30. 2
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java
  31. 139
      common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java
  32. 10
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  33. 96
      dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java
  34. 69
      dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java
  35. 72
      dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java
  36. 86
      dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java
  37. 33
      dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java
  38. 11
      dao/src/main/resources/sql/schema-entities.sql

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

@ -0,0 +1,83 @@
/**
* 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.EntityType;
import org.thingsboard.server.common.data.ProfileEntityIdInfo;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.job.CfReprocessingJobConfiguration;
import org.thingsboard.server.common.data.job.CfReprocessingTask;
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 org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.device.DeviceService;
import java.util.function.Consumer;
@Component
@RequiredArgsConstructor
public class CfReprocessingJobProcessor extends JobProcessor {
private final DeviceService deviceService;
private final AssetService assetService;
@Override
public void process(Job job, Consumer<Task> taskConsumer) {
CfReprocessingJobConfiguration configuration = job.getConfiguration();
CalculatedField calculatedField = configuration.getCalculatedField();
EntityId entityId = calculatedField.getEntityId();
if (entityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) {
taskConsumer.accept(createTask(job, configuration, entityId));
} else {
PageDataIterable<ProfileEntityIdInfo> entities;
if (entityId.getEntityType() == EntityType.DEVICE_PROFILE) {
entities = new PageDataIterable<>(pageLink -> deviceService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512);
} else if (entityId.getEntityType() == EntityType.ASSET_PROFILE) {
entities = new PageDataIterable<>(pageLink -> assetService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512);
} else {
throw new IllegalArgumentException("Unsupported CF entity type " + entityId.getEntityType());
}
entities.forEach(device -> {
taskConsumer.accept(createTask(job, configuration, device.getEntityId()));
});
}
}
private Task createTask(Job job, CfReprocessingJobConfiguration configuration, EntityId entityId) {
return CfReprocessingTask.builder()
.tenantId(job.getTenantId())
.jobId(job.getId())
.key(entityId.getEntityType().getNormalName() + " " + entityId.getId())
.calculatedField(configuration.getCalculatedField())
.entityId(entityId)
.startTs(configuration.getStartTs())
.endTs(configuration.getEndTs())
.build();
}
@Override
public JobType getType() {
return JobType.CF_REPROCESSING;
}
}

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

@ -0,0 +1,136 @@
/**
* 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 jakarta.annotation.PreDestroy;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.id.JobId;
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 org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.task.JobService;
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.TbQueueMsgMetadata;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.QueueConsumerManager;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.function.Function;
import java.util.stream.Collectors;
@TbCoreComponent
@Component
@Slf4j
public class DefaultJobManager implements JobManager {
private final JobService jobService;
private final TbCoreQueueFactory queueFactory;
private final Map<JobType, JobProcessor> jobProcessors;
private final Map<JobType, TbQueueProducer<TbProtoQueueMsg<TaskProto>>> taskProducers;
private final QueueConsumerManager<TbProtoQueueMsg<TaskResultProto>> taskResultConsumer;
private final ExecutorService consumerExecutor;
public DefaultJobManager(JobService jobService, TbCoreQueueFactory queueFactory, List<JobProcessor> jobProcessors) {
this.jobService = jobService;
this.queueFactory = queueFactory;
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.consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("task-result-consumer"));
this.taskResultConsumer = QueueConsumerManager.<TbProtoQueueMsg<TaskResultProto>>builder() // fixme: should be consumer per partition
.name("tasks-results")
.msgPackProcessor(this::processResults)
.pollInterval(125)
.consumerCreator(queueFactory::createTaskResultConsumer)
.consumerExecutor(consumerExecutor)
.build();
}
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
public void afterStartUp() {
taskResultConsumer.subscribe();
taskResultConsumer.launch();
}
@Override
public void submitJob(Job job) {
job = jobService.createJob(job.getTenantId(), job);
log.info("Submitting job: {}", job);
jobProcessors.get(job.getType()).process(job, this::submitTask);
}
private void submitTask(Task task) {
log.info("Submitting task: {}", task);
TaskProto taskProto = TaskProto.newBuilder()
.setValue(JacksonUtil.toString(task))
.build();
TbQueueProducer<TbProtoQueueMsg<TaskProto>> producer = taskProducers.get(task.getJobType());
TbProtoQueueMsg<TaskProto> msg = new TbProtoQueueMsg<>(task.getTenantId().getId(), taskProto); // one job at a time for a given tenant
producer.send(TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(), msg, new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
log.trace("Submitted task: {}", task);
}
@Override
public void onFailure(Throwable t) {
log.warn("Failed to submit task: {}", task, t);
}
});
}
@SneakyThrows
private void processResults(List<TbProtoQueueMsg<TaskResultProto>> msgs, TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> consumer) {
Map<JobId, List<TaskResult>> results = msgs.stream()
.map(msg -> JacksonUtil.fromString(msg.getValue().getValue(), TaskResult.class))
.collect(Collectors.groupingBy(TaskResult::getJobId));
results.forEach((jobId, taskResults) -> {
try {
log.info("[{}] Processing task results: {}", jobId, taskResults);
jobService.reportTaskResults(jobId, taskResults);
} catch (Exception e) {
log.warn("Failed to report task results for job {}: {}", jobId, taskResults, e);
}
});
consumer.commit();
Thread.sleep(5000);
}
@PreDestroy
private void destroy() {
taskResultConsumer.stop();
consumerExecutor.shutdownNow();
}
}

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

@ -0,0 +1,24 @@
/**
* 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.thingsboard.server.common.data.job.Job;
public interface JobManager {
void submitJob(Job job);
}

30
application/src/main/java/org/thingsboard/server/service/job/JobProcessor.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.service.job;
import org.thingsboard.server.common.data.job.Task;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobType;
import java.util.function.Consumer;
public abstract class JobProcessor {
public abstract void process(Job job, Consumer<Task> taskConsumer);
public abstract JobType getType();
}

57
application/src/main/java/org/thingsboard/server/service/job/task/CfReprocessingTaskProcessor.java

@ -0,0 +1,57 @@
/**
* 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 com.google.common.util.concurrent.SettableFuture;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
import org.thingsboard.server.actors.calculatedField.CalculatedFieldReprocessingService;
import org.thingsboard.server.common.data.job.CfReprocessingTask;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.queue.task.TaskProcessor;
import java.util.concurrent.TimeUnit;
@Component
@RequiredArgsConstructor
public class CfReprocessingTaskProcessor extends TaskProcessor<CfReprocessingTask> {
private final CalculatedFieldReprocessingService cfReprocessingService;
@Override
protected void process(CfReprocessingTask task) throws Exception {
SettableFuture<Void> future = SettableFuture.create();
cfReprocessingService.reprocess(task, new TbCallback() {
@Override
public void onSuccess() {
future.set(null);
}
@Override
public void onFailure(Throwable t) {
future.setException(t);
}
});
future.get(1, TimeUnit.MINUTES);
}
@Override
public JobType getJobType() {
return JobType.CF_REPROCESSING;
}
}

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

@ -1667,6 +1667,8 @@ queue:
edqs-requests: "${TB_QUEUE_KAFKA_EDQS_REQUESTS_TOPIC_PROPERTIES:retention.ms:180000;segment.bytes:52428800;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" edqs-requests: "${TB_QUEUE_KAFKA_EDQS_REQUESTS_TOPIC_PROPERTIES:retention.ms:180000;segment.bytes:52428800;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}"
# Kafka properties for EDQS state topic (infinite retention, compaction) # Kafka properties for EDQS state topic (infinite retention, compaction)
edqs-state: "${TB_QUEUE_KAFKA_EDQS_STATE_TOPIC_PROPERTIES:retention.ms:-1;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1;cleanup.policy:compact}" edqs-state: "${TB_QUEUE_KAFKA_EDQS_STATE_TOPIC_PROPERTIES:retention.ms:-1;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1;cleanup.policy:compact}"
# Kafka properties for tasks topics
tasks: "${TB_QUEUE_KAFKA_TASKS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:52428800;retention.bytes:104857600;partitions:100;min.insync.replicas:1}"
consumer-stats: consumer-stats:
# Prints lag between consumer group offset and last messages offset in Kafka topics # Prints lag between consumer group offset and last messages offset in Kafka topics
enabled: "${TB_QUEUE_KAFKA_CONSUMER_STATS_ENABLED:true}" enabled: "${TB_QUEUE_KAFKA_CONSUMER_STATS_ENABLED:true}"

2
application/src/test/resources/logback-test.xml

@ -9,7 +9,7 @@
<!-- <logger name="org.thingsboard.server.service.subscription" level="TRACE"/>--> <!-- <logger name="org.thingsboard.server.service.subscription" level="TRACE"/>-->
<logger name="org.thingsboard.server.controller.TbTestWebSocketClient" level="INFO"/> <logger name="org.thingsboard.server.controller.TbTestWebSocketClient" level="INFO"/>
<logger name="org.thingsboard.server" level="WARN"/> <logger name="org.thingsboard.server" level="INFO"/>
<logger name="org.springframework" level="WARN"/> <logger name="org.springframework" level="WARN"/>
<logger name="org.springframework.boot.test" level="WARN"/> <logger name="org.springframework.boot.test" level="WARN"/>
<logger name="org.apache.cassandra" level="WARN"/> <logger name="org.apache.cassandra" level="WARN"/>

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

@ -0,0 +1,35 @@
/**
* 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.dao.task;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.job.Job;
import java.util.List;
public interface JobService {
Job createJob(TenantId tenantId, Job job);
void reportTaskResults(JobId jobId, List<TaskResult> results);
PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink);
}

3
common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java

@ -63,7 +63,8 @@ public enum EntityType {
MOBILE_APP(37), MOBILE_APP(37),
MOBILE_APP_BUNDLE(38), MOBILE_APP_BUNDLE(38),
CALCULATED_FIELD(39), CALCULATED_FIELD(39),
CALCULATED_FIELD_LINK(40); CALCULATED_FIELD_LINK(40),
JOB(41);
@Getter @Getter
private final int protoNumber; // Corresponds to EntityTypeProto private final int protoNumber; // Corresponds to EntityTypeProto

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

@ -117,6 +117,8 @@ public class EntityIdFactory {
return new CalculatedFieldId(uuid); return new CalculatedFieldId(uuid);
case CALCULATED_FIELD_LINK: case CALCULATED_FIELD_LINK:
return new CalculatedFieldLinkId(uuid); return new CalculatedFieldLinkId(uuid);
case JOB:
return new JobId(uuid);
} }
throw new IllegalArgumentException("EntityType " + type + " is not supported!"); throw new IllegalArgumentException("EntityType " + type + " is not supported!");
} }

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

@ -0,0 +1,38 @@
/**
* 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.id;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.swagger.v3.oas.annotations.media.Schema;
import org.thingsboard.server.common.data.EntityType;
import java.util.UUID;
public class JobId extends UUIDBased implements EntityId {
@JsonCreator
public JobId(@JsonProperty("id") UUID id) {
super(id);
}
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "string", example = "TASK", allowableValues = "TASK")
@Override
public EntityType getEntityType() {
return EntityType.JOB;
}
}

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

@ -0,0 +1,39 @@
/**
* 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 org.thingsboard.server.common.data.cf.CalculatedField;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class CfReprocessingJobConfiguration implements JobConfiguration {
private CalculatedField calculatedField;
private long startTs;
private long endTs;
@Override
public JobType getType() {
return JobType.CF_REPROCESSING;
}
}

25
common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobResult.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 CfReprocessingJobResult extends JobResult {
@Override
public JobType getJobType() {
return JobType.CF_REPROCESSING;
}
}

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

@ -0,0 +1,53 @@
/**
* 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.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.cf.CalculatedField;
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
@AllArgsConstructor
@NoArgsConstructor
@EqualsAndHashCode(callSuper = true)
public class CfReprocessingTask extends Task {
private CalculatedField calculatedField;
private EntityId entityId;
private long startTs;
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
public JobType getJobType() {
return JobType.CF_REPROCESSING;
}
}

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

@ -0,0 +1,56 @@
/**
* 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 lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
@Data
@NoArgsConstructor
@EqualsAndHashCode(callSuper = true)
public class Job extends BaseData<JobId> implements HasTenantId {
private TenantId tenantId;
private JobType type;
private String key;
private JobStatus status;
private JobConfiguration configuration;
private JobResult result;
@Builder
public Job(TenantId tenantId, JobType type, String key, JobConfiguration configuration) {
this.tenantId = tenantId;
this.type = type;
this.key = key;
this.configuration = configuration;
this.status = JobStatus.PENDING;
this.result = switch (type) {
case CF_REPROCESSING -> new CfReprocessingJobResult();
};
}
@SuppressWarnings("unchecked")
public <C extends JobConfiguration> C getConfiguration() {
return (C) configuration;
}
}

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

@ -0,0 +1,32 @@
/**
* 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type")
@JsonSubTypes({
@Type(name = "CF_REPROCESSING", value = CfReprocessingJobConfiguration.class),
})
public interface JobConfiguration {
JobType getType();
}

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

@ -0,0 +1,41 @@
/**
* 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.Data;
import java.util.HashMap;
import java.util.Map;
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType")
@JsonSubTypes({
@JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class),
})
@Data
public abstract class JobResult {
private int successfulCount;
private int failedCount;
private int totalCount;
private Map<String, String> failures = new HashMap<>();
public abstract JobType getJobType();
}

24
common/data/src/main/java/org/thingsboard/server/common/data/job/JobStatus.java

@ -0,0 +1,24 @@
/**
* 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 enum JobStatus {
PENDING,
RUNNING,
COMPLETED,
FAILED,
CANCELLED
}

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

@ -0,0 +1,26 @@
/**
* 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 enum JobType {
CF_REPROCESSING;
public String getTasksTopic() {
return "tasks." + name().toLowerCase();
}
}

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

@ -0,0 +1,50 @@
/**
* 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import lombok.Data;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
@Data
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType")
@JsonSubTypes({
@JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class),
})
public abstract class Task {
private TenantId tenantId;
private JobId jobId;
private String key;
public Task(TenantId tenantId, JobId jobId, String key) {
this.tenantId = tenantId;
this.jobId = jobId;
this.key = key;
}
public Task() {
}
private int attempt = 0;
public abstract JobType getJobType();
}

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

@ -0,0 +1,45 @@
/**
* 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 org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.id.TenantId;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class TaskResult {
private TenantId tenantId;
private JobId jobId;
private boolean success;
private TaskFailure failure;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public static class TaskFailure {
private String error;
private Task task;
}
}

8
common/proto/src/main/proto/queue.proto

@ -1846,3 +1846,11 @@ message EdqsRequestMsg {
message EdqsResponseMsg { message EdqsResponseMsg {
string value = 1; string value = 1;
} }
message TaskProto {
string value = 1; // fixme: TMP, make more efficient
}
message TaskResultProto {
string value = 1; // fixme: TMP, make more efficient
}

5
common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java

@ -62,6 +62,8 @@ public class TbKafkaTopicConfigs {
private String edqsRequestsProperties; private String edqsRequestsProperties;
@Value("${queue.kafka.topic-properties.edqs-state:}") @Value("${queue.kafka.topic-properties.edqs-state:}")
private String edqsStateProperties; private String edqsStateProperties;
@Value("${queue.kafka.topic-properties.tasks:}")
private String tasksProperties;
@Getter @Getter
private Map<String, String> coreConfigs; private Map<String, String> coreConfigs;
@ -99,6 +101,8 @@ public class TbKafkaTopicConfigs {
private Map<String, String> edqsRequestsConfigs; private Map<String, String> edqsRequestsConfigs;
@Getter @Getter
private Map<String, String> edqsStateConfigs; private Map<String, String> edqsStateConfigs;
@Getter
private Map<String, String> tasksConfigs;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -122,6 +126,7 @@ public class TbKafkaTopicConfigs {
edqsEventsConfigs = PropertyUtils.getProps(edqsEventsProperties); edqsEventsConfigs = PropertyUtils.getProps(edqsEventsProperties);
edqsRequestsConfigs = PropertyUtils.getProps(edqsRequestsProperties); edqsRequestsConfigs = PropertyUtils.getProps(edqsRequestsProperties);
edqsStateConfigs = PropertyUtils.getProps(edqsStateProperties); edqsStateConfigs = PropertyUtils.getProps(edqsStateProperties);
tasksConfigs = PropertyUtils.getProps(tasksProperties);
} }
} }

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

@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.queue.Queue; 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;
@ -258,9 +259,28 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
.build(); .build();
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.TaskProto>> createTaskProducer(JobType jobType) {
return new InMemoryTbQueueProducer<>(storage, jobType.getTasksTopic());
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.TaskProto>> createTaskConsumer(JobType jobType) {
return new InMemoryTbQueueConsumer<>(storage, jobType.getTasksTopic());
}
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultProducer() {
return new InMemoryTbQueueProducer<>(storage, "tasks.results");
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultConsumer() {
return new InMemoryTbQueueConsumer<>(storage, "tasks.results");
}
@Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}") @Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}")
private void printInMemoryStats() { private void printInMemoryStats() {
storage.printStats(); storage.printStats();
} }
} }

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

@ -23,12 +23,15 @@ import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.EdgeId; 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.queue.Queue; 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.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.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;
@ -110,6 +113,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin cfStateAdmin;
private final TbQueueAdmin edqsEventsAdmin; private final TbQueueAdmin edqsEventsAdmin;
private final TbKafkaAdmin edqsRequestsAdmin; private final TbKafkaAdmin edqsRequestsAdmin;
private final TbQueueAdmin tasksAdmin;
private final AtomicLong consumerCount = new AtomicLong(); private final AtomicLong consumerCount = new AtomicLong();
private final AtomicLong edgeConsumerCount = new AtomicLong(); private final AtomicLong edgeConsumerCount = new AtomicLong();
@ -158,6 +162,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs()); this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs());
this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs());
this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs()); this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs());
this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs());
} }
@Override @Override
@ -641,6 +646,52 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
.build(); .build();
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TaskProto>> createTaskProducer(JobType jobType) {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TaskProto>>builder()
.clientId(jobType.name().toLowerCase() + "-task-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName(jobType.getTasksTopic()))
.settings(kafkaSettings)
.admin(tasksAdmin)
.build();
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType) {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TaskProto>>builder()
.settings(kafkaSettings)
.topic(topicService.buildTopicName(jobType.getTasksTopic()))
.clientId(jobType.name().toLowerCase() + "-task-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName(jobType.name().toLowerCase() + "-task-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskProto.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin)
.statsService(consumerStatsService)
.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> createTaskResultProducer() {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TaskResultProto>>builder()
.clientId("task-result-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName("tasks.results"))
.settings(kafkaSettings)
.admin(tasksAdmin)
.build();
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> createTaskResultConsumer() {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TaskResultProto>>builder()
.settings(kafkaSettings)
.topic(topicService.buildTopicName("tasks.results"))
.clientId("task-result-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName("task-result-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskResultProto.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin)
.statsService(consumerStatsService)
.build();
}
@PreDestroy @PreDestroy
private void destroy() { private void destroy() {
if (coreAdmin != null) { if (coreAdmin != null) {

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

@ -22,9 +22,12 @@ import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EdgeId; 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.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.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;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -105,6 +108,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfAdmin;
private final TbQueueAdmin edqsEventsAdmin; private final TbQueueAdmin edqsEventsAdmin;
private final TbKafkaAdmin edqsRequestsAdmin; private final TbKafkaAdmin edqsRequestsAdmin;
private final TbQueueAdmin tasksAdmin;
private final AtomicLong consumerCount = new AtomicLong(); private final AtomicLong consumerCount = new AtomicLong();
private final AtomicLong edgeConsumerCount = new AtomicLong(); private final AtomicLong edgeConsumerCount = new AtomicLong();
@ -153,6 +157,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs()); this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs());
this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs());
this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs()); this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs());
this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs());
} }
@Override @Override
@ -520,6 +525,29 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
.build(); .build();
} }
@Override
public TbQueueProducer<TbProtoQueueMsg<TaskProto>> createTaskProducer(JobType jobType) {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TaskProto>>builder()
.clientId(jobType.name().toLowerCase() + "-task-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName(jobType.getTasksTopic()))
.settings(kafkaSettings)
.admin(tasksAdmin)
.build();
}
@Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultConsumer() {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TransportProtos.TaskResultProto>>builder()
.settings(kafkaSettings)
.topic(topicService.buildTopicName("tasks.results"))
.clientId("task-result-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName("task-result-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.TaskResultProto.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin)
.statsService(consumerStatsService)
.build();
}
@PreDestroy @PreDestroy
private void destroy() { private void destroy() {
if (coreAdmin != null) { if (coreAdmin != null) {

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

@ -22,12 +22,15 @@ import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
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.queue.Queue; 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.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;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
@ -96,6 +99,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfAdmin;
private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin cfStateAdmin;
private final TbQueueAdmin edqsEventsAdmin; private final TbQueueAdmin edqsEventsAdmin;
private final TbQueueAdmin tasksAdmin;
private final AtomicLong consumerCount = new AtomicLong(); private final AtomicLong consumerCount = new AtomicLong();
public KafkaTbRuleEngineQueueFactory(TopicService topicService, TbKafkaSettings kafkaSettings, public KafkaTbRuleEngineQueueFactory(TopicService topicService, TbKafkaSettings kafkaSettings,
@ -133,6 +137,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs()); this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs());
this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs()); this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs());
this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs());
this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs());
} }
@Override @Override
@ -414,6 +419,29 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
throw new UnsupportedOperationException(); throw new UnsupportedOperationException();
} }
@Override
public TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType) {
return TbKafkaConsumerTemplate.<TbProtoQueueMsg<TaskProto>>builder()
.settings(kafkaSettings)
.topic(topicService.buildTopicName(jobType.getTasksTopic()))
.clientId(jobType.name().toLowerCase() + "-task-consumer-" + serviceInfoProvider.getServiceId())
.groupId(topicService.buildTopicName(jobType.name().toLowerCase() + "-task-consumer-group"))
.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskProto.parseFrom(msg.getData()), msg.getHeaders()))
.admin(tasksAdmin)
.statsService(consumerStatsService)
.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<TransportProtos.TaskResultProto>> createTaskResultProducer() {
return TbKafkaProducerTemplate.<TbProtoQueueMsg<TransportProtos.TaskResultProto>>builder()
.clientId("task-result-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName("tasks.results"))
.settings(kafkaSettings)
.admin(tasksAdmin)
.build();
}
@PreDestroy @PreDestroy
private void destroy() { private void destroy() {
if (coreAdmin != null) { if (coreAdmin != null) {
@ -438,4 +466,5 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
cfAdmin.destroy(); cfAdmin.destroy();
} }
} }
} }

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

@ -0,0 +1,31 @@
/**
* 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.provider;
import org.thingsboard.server.common.data.job.JobType;
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.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
public interface TaskProcessorQueueFactory {
TbQueueConsumer<TbProtoQueueMsg<TaskProto>> createTaskConsumer(JobType jobType);
TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> createTaskResultProducer();
}

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

@ -17,7 +17,10 @@ package org.thingsboard.server.queue.provider;
import org.thingsboard.server.common.data.id.EdgeId; 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.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
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;
@ -165,4 +168,8 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer();
TbQueueProducer<TbProtoQueueMsg<TaskProto>> createTaskProducer(JobType jobType);
TbQueueConsumer<TbProtoQueueMsg<TaskResultProto>> createTaskResultConsumer();
} }

2
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java

@ -41,7 +41,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 TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory { public interface TbRuleEngineQueueFactory 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

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

@ -0,0 +1,139 @@
/**
* 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 jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.common.data.job.Task;
import org.thingsboard.server.common.data.job.TaskResult;
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.TaskResultProto;
import org.thingsboard.server.queue.TbQueueCallback;
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.consumer.QueueConsumerManager;
import org.thingsboard.server.queue.provider.TaskProcessorQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@Slf4j
public abstract class TaskProcessor<T extends Task> {
@Autowired
private TaskProcessorQueueFactory queueFactory;
private QueueConsumerManager<TbProtoQueueMsg<TaskProto>> taskConsumer;
private TbQueueProducer<TbProtoQueueMsg<TaskResultProto>> taskResultProducer;
private ExecutorService consumerExecutor;
@PostConstruct
public void init() {
consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-consumer"));
taskConsumer = QueueConsumerManager.<TbProtoQueueMsg<TaskProto>>builder() // fixme: should be consumer per partition
.name(getJobType().name() + "-tasks")
.msgPackProcessor(this::processMsgs)
.pollInterval(125)
.consumerCreator(() -> queueFactory.createTaskConsumer(getJobType()))
.consumerExecutor(consumerExecutor)
.build();
taskResultProducer = queueFactory.createTaskResultProducer();
}
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
public void afterStartUp() {
taskConsumer.subscribe();
taskConsumer.launch();
}
@PreDestroy
public void destroy() {
taskConsumer.stop();
consumerExecutor.shutdownNow();
}
private void processMsgs(List<TbProtoQueueMsg<TaskProto>> msgs, TbQueueConsumer<TbProtoQueueMsg<TaskProto>> consumer) {
for (TbProtoQueueMsg<TaskProto> msg : msgs) {
TaskProto taskProto = msg.getValue();
Task task = JacksonUtil.fromString(taskProto.getValue(), Task.class);
processTask((T) task);
}
consumer.commit();
}
private void processTask(T task) {
task.setAttempt(task.getAttempt() + 1);
log.info("Processing task: {}", task);
try {
process(task);
reportSuccess(task);
} catch (Exception e) {
log.error("Failed to process task (attempt {}): {}", task.getAttempt(), task, e);
if (task.getAttempt() < 3) {
processTask(task);
} else {
reportFailure(task, e);
}
}
}
private void reportSuccess(Task task) {
TaskResult result = TaskResult.builder()
.tenantId(task.getTenantId())
.jobId(task.getJobId())
.success(true)
.build();
reportResult(result);
}
private void reportFailure(Task task, Throwable error) {
TaskResult result = TaskResult.builder()
.tenantId(task.getTenantId())
.jobId(task.getJobId())
.failure(TaskFailure.builder()
.error(error.getMessage())
.task(task)
.build())
.build();
reportResult(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;
public abstract JobType getJobType();
}

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

@ -739,6 +739,16 @@ public class ModelConstants {
public static final String CALCULATED_FIELD_LINK_ENTITY_ID = ENTITY_ID_COLUMN; public static final String CALCULATED_FIELD_LINK_ENTITY_ID = ENTITY_ID_COLUMN;
public static final String CALCULATED_FIELD_LINK_CALCULATED_FIELD_ID = "calculated_field_id"; public static final String CALCULATED_FIELD_LINK_CALCULATED_FIELD_ID = "calculated_field_id";
/**
* Tasks constants.
*/
public static final String JOB_TABLE_NAME = "job";
public static final String JOB_TYPE_PROPERTY = "type";
public static final String JOB_KEY_PROPERTY = "key";
public static final String JOB_STATUS_PROPERTY = "status";
public static final String JOB_CONFIGURATION_PROPERTY = "configuration";
public static final String JOB_RESULT_PROPERTY = "result";
protected static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, JSON_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN}; protected static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, JSON_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN};
protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN), max(TS_COLUMN)}; protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN), max(TS_COLUMN)};

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

@ -0,0 +1,96 @@
/**
* 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.dao.model.sql;
import com.fasterxml.jackson.databind.JsonNode;
import jakarta.persistence.Column;
import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
import jakarta.persistence.EnumType;
import jakarta.persistence.Enumerated;
import jakarta.persistence.Table;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import org.hibernate.annotations.JdbcType;
import org.hibernate.dialect.PostgreSQLJsonPGObjectJsonbType;
import org.thingsboard.server.common.data.id.JobId;
import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.job.JobConfiguration;
import org.thingsboard.server.common.data.job.JobResult;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.JobType;
import org.thingsboard.server.dao.model.BaseSqlEntity;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.util.mapping.JsonConverter;
import java.util.UUID;
@Data
@EqualsAndHashCode(callSuper = true)
@NoArgsConstructor
@Entity
@Table(name = ModelConstants.JOB_TABLE_NAME)
public class JobEntity extends BaseSqlEntity<Job> {
@Column(name = ModelConstants.TENANT_ID_PROPERTY, nullable = false)
private UUID tenantId;
@Enumerated(EnumType.STRING)
@Column(name = ModelConstants.JOB_TYPE_PROPERTY, nullable = false)
private JobType type;
@Column(name = ModelConstants.JOB_KEY_PROPERTY, nullable = false)
private String key;
@Enumerated(EnumType.STRING)
@Column(name = ModelConstants.JOB_STATUS_PROPERTY, nullable = false)
private JobStatus status;
@Convert(converter = JsonConverter.class)
@Column(name = ModelConstants.JOB_CONFIGURATION_PROPERTY, nullable = false)
private JsonNode configuration;
@Convert(converter = JsonConverter.class)
@JdbcType(PostgreSQLJsonPGObjectJsonbType.class)
@Column(name = ModelConstants.JOB_RESULT_PROPERTY)
private JsonNode result;
public JobEntity(Job job) {
super(job);
this.tenantId = getTenantUuid(job.getTenantId());
this.type = job.getType();
this.key = job.getKey();
this.status = job.getStatus();
this.configuration = toJson(job.getConfiguration());
this.result = toJson(job.getResult());
}
@Override
public Job toData() {
Job job = new Job();
job.setId(new JobId(id));
job.setCreatedTime(createdTime);
job.setTenantId(getTenantId(tenantId));
job.setType(type);
job.setKey(key);
job.setStatus(status);
job.setConfiguration(fromJson(configuration, JobConfiguration.class));
job.setResult(fromJson(result, JobResult.class));
return job;
}
}

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

@ -0,0 +1,69 @@
/**
* 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.dao.sql.task;
import jakarta.transaction.Transactional;
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.Modifying;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
import org.springframework.stereotype.Repository;
import org.thingsboard.server.dao.model.sql.JobEntity;
import java.util.UUID;
@Repository
public interface JobRepository extends JpaRepository<JobEntity, UUID> {
Page<JobEntity> findByTenantId(UUID tenantId, Pageable pageable);
@Modifying
@Transactional
@Query(value = """
UPDATE job
SET result = jsonb_set(
result,
'{successfulCount}',
to_jsonb((result->>'successfulCount')::int + :count)
)
WHERE id = :jobId
RETURNING ((result->>'successfulCount')::int + :count)
+ (result->>'failedCount')::int = (result->>'totalCount')::int
""", nativeQuery = true)
boolean reportTaskSuccess(@Param("jobId") UUID jobId, @Param("count") int count);
@Modifying
@Transactional
@Query(value = """
UPDATE job
SET result = jsonb_set(
jsonb_set(
result,
'{failedCount}',
to_jsonb((result->>'failedCount')::int + 1)
),
ARRAY['failures', :taskKey],
to_jsonb(:error)
)
WHERE id = :jobId
RETURNING ((result->>'failedCount')::int + 1) + (result->>'successfulCount')::int
= (result->>'totalCount')::int
""", nativeQuery = true)
boolean reportTaskFailure(@Param("jobId") UUID jobId, @Param("taskKey") String taskKey, @Param("error") String error);
}

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

@ -0,0 +1,72 @@
/**
* 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.dao.sql.task;
import lombok.RequiredArgsConstructor;
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.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sql.JobEntity;
import org.thingsboard.server.dao.sql.JpaAbstractDao;
import org.thingsboard.server.dao.task.JobDao;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.UUID;
@Component
@SqlDao
@RequiredArgsConstructor
public class JpaJobDao extends JpaAbstractDao<JobEntity, Job> implements JobDao {
private final JobRepository jobRepository;
@Override
public PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink) {
return DaoUtil.toPageData(jobRepository.findByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink)));
}
@Override
public boolean reportTaskSuccess(JobId jobId, int tasksCount) {
return jobRepository.reportTaskSuccess(jobId.getId(), tasksCount);
}
@Override
public boolean reportTaskFailure(JobId jobId, String taskKey, String error) {
return jobRepository.reportTaskFailure(jobId.getId(), taskKey, error);
}
@Override
public EntityType getEntityType() {
return EntityType.JOB;
}
@Override
protected Class<JobEntity> getEntityClass() {
return JobEntity.class;
}
@Override
protected JpaRepository<JobEntity, UUID> getRepository() {
return jobRepository;
}
}

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

@ -0,0 +1,86 @@
/**
* 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.dao.task;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
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.JobResult;
import org.thingsboard.server.common.data.job.JobStatus;
import org.thingsboard.server.common.data.job.TaskResult;
import org.thingsboard.server.common.data.job.TaskResult.TaskFailure;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import java.util.List;
@Service
@RequiredArgsConstructor
@Slf4j
public class DefaultJobService implements JobService {
private final JobDao jobDao;
@Override
public Job createJob(TenantId tenantId, Job job) {
return jobDao.save(tenantId, job);
}
@Override
public void reportTaskResults(JobId jobId, List<TaskResult> results) {
Job job = jobDao.findById(TenantId.SYS_TENANT_ID, jobId.getId());
switch (job.getStatus()) {
case PENDING -> {
job.setStatus(JobStatus.RUNNING);
}
case CANCELLED, COMPLETED, FAILED -> {
// got some stale stats
return;
}
}
JobResult jobResult = job.getResult();
for (TaskResult taskResult : results) {
if (taskResult.isSuccess()) {
jobResult.setSuccessfulCount(jobResult.getSuccessfulCount() + 1);
} else {
TaskFailure failure = taskResult.getFailure();
String key = failure.getTask().getKey();
jobResult.setFailedCount(jobResult.getFailedCount() + 1);
jobResult.getFailures().put(key, failure.getError());
}
}
if (jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) {
if (jobResult.getFailures().isEmpty()) {
job.setStatus(JobStatus.COMPLETED);
} else {
job.setStatus(JobStatus.FAILED);
}
}
log.info("Saving job {}", job);
jobDao.save(TenantId.SYS_TENANT_ID, job);
}
@Override
public PageData<Job> findJobsByTenantId(TenantId tenantId, PageLink pageLink) {
return jobDao.findByTenantId(tenantId, pageLink);
}
}

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

@ -0,0 +1,33 @@
/**
* 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.dao.task;
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.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.Dao;
public interface JobDao extends Dao<Job> {
PageData<Job> findByTenantId(TenantId tenantId, PageLink pageLink);
boolean reportTaskSuccess(JobId jobId, int tasksCount);
boolean reportTaskFailure(JobId jobId, String taskKey, String error);
}

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

@ -948,3 +948,14 @@ CREATE TABLE IF NOT EXISTS cf_debug_event (
e_result varchar, e_result varchar,
e_error varchar e_error varchar
) PARTITION BY RANGE (ts); ) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS job (
id uuid NOT NULL CONSTRAINT job_pkey PRIMARY KEY,
created_time bigint NOT NULL,
tenant_id uuid NOT NULL,
type varchar NOT NULL,
key varchar NOT NULL,
status varchar NOT NULL,
configuration varchar(1000) NOT NULL,
result jsonb
);

Loading…
Cancel
Save