|
|
|
@ -18,12 +18,12 @@ package org.thingsboard.server.service.job; |
|
|
|
import jakarta.annotation.PreDestroy; |
|
|
|
import lombok.SneakyThrows; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.apache.commons.lang3.exception.ExceptionUtils; |
|
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
|
import org.springframework.stereotype.Component; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
|
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
|
|
import org.thingsboard.rule.engine.api.NotificationCenter; |
|
|
|
import org.thingsboard.server.common.data.id.JobId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.job.Job; |
|
|
|
@ -33,8 +33,12 @@ import org.thingsboard.server.common.data.job.JobStatus; |
|
|
|
import org.thingsboard.server.common.data.job.JobType; |
|
|
|
import org.thingsboard.server.common.data.job.task.Task; |
|
|
|
import org.thingsboard.server.common.data.job.task.TaskResult; |
|
|
|
import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; |
|
|
|
import org.thingsboard.server.common.data.notification.targets.platform.TenantAdministratorsFilter; |
|
|
|
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; |
|
|
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
|
|
|
import org.thingsboard.server.dao.job.JobService; |
|
|
|
import org.thingsboard.server.dao.notification.DefaultNotifications; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; |
|
|
|
import org.thingsboard.server.queue.TbQueueCallback; |
|
|
|
@ -65,6 +69,7 @@ public class DefaultJobManager implements JobManager { |
|
|
|
|
|
|
|
private final JobService jobService; |
|
|
|
private final JobStatsService jobStatsService; |
|
|
|
private final NotificationCenter notificationCenter; |
|
|
|
private final Map<JobType, JobProcessor> jobProcessors; |
|
|
|
private final Map<JobType, TbQueueProducer<TbProtoQueueMsg<TaskProto>>> taskProducers; |
|
|
|
private final QueueConsumerManager<TbProtoQueueMsg<JobStatsMsg>> jobStatsConsumer; |
|
|
|
@ -74,9 +79,11 @@ public class DefaultJobManager implements JobManager { |
|
|
|
@Value("${queue.tasks.stats.processing_interval_ms:5000}") |
|
|
|
private int statsProcessingInterval; |
|
|
|
|
|
|
|
public DefaultJobManager(JobService jobService, JobStatsService jobStatsService, TbCoreQueueFactory queueFactory, List<JobProcessor> jobProcessors) { |
|
|
|
public DefaultJobManager(JobService jobService, JobStatsService jobStatsService, NotificationCenter notificationCenter, |
|
|
|
TbCoreQueueFactory queueFactory, List<JobProcessor> jobProcessors) { |
|
|
|
this.jobService = jobService; |
|
|
|
this.jobStatsService = jobStatsService; |
|
|
|
this.notificationCenter = notificationCenter; |
|
|
|
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.executor = ThingsBoardExecutors.newWorkStealingPool(Math.max(4, Runtime.getRuntime().availableProcessors()), getClass()); |
|
|
|
@ -104,10 +111,29 @@ public class DefaultJobManager implements JobManager { |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onJobUpdate(Job job) { |
|
|
|
if (job.getStatus() == JobStatus.PENDING) { |
|
|
|
executor.execute(() -> { |
|
|
|
processJob(job); |
|
|
|
}); |
|
|
|
JobStatus status = job.getStatus(); |
|
|
|
switch (status) { |
|
|
|
case PENDING -> { |
|
|
|
executor.execute(() -> { |
|
|
|
try { |
|
|
|
processJob(job); |
|
|
|
} catch (Throwable e) { |
|
|
|
log.error("Failed to process job update: {}", job, e); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
case COMPLETED, FAILED -> { |
|
|
|
executor.execute(() -> { |
|
|
|
try { |
|
|
|
if (status == JobStatus.COMPLETED) { |
|
|
|
getJobProcessor(job.getType()).onJobCompleted(job); |
|
|
|
} |
|
|
|
sendJobFinishedNotification(job); |
|
|
|
} catch (Throwable e) { |
|
|
|
log.error("Failed to process job update: {}", job, e); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -115,7 +141,7 @@ public class DefaultJobManager implements JobManager { |
|
|
|
TenantId tenantId = job.getTenantId(); |
|
|
|
JobId jobId = job.getId(); |
|
|
|
try { |
|
|
|
JobProcessor processor = jobProcessors.get(job.getType()); |
|
|
|
JobProcessor processor = getJobProcessor(job.getType()); |
|
|
|
List<TaskResult> toReprocess = job.getConfiguration().getToReprocess(); |
|
|
|
if (toReprocess == null) { |
|
|
|
int tasksCount = processor.process(job, this::submitTask); // todo: think about stopping tb - while tasks are being submitted
|
|
|
|
@ -127,11 +153,7 @@ public class DefaultJobManager implements JobManager { |
|
|
|
} |
|
|
|
} catch (Throwable e) { |
|
|
|
log.error("[{}][{}][{}] Failed to submit tasks", tenantId, jobId, job.getType(), e); |
|
|
|
try { |
|
|
|
jobService.markAsFailed(tenantId, jobId, ExceptionUtils.getStackTrace(e)); |
|
|
|
} catch (Throwable e2) { |
|
|
|
log.error("[{}][{}] Failed to mark job as failed", tenantId, jobId, e2); |
|
|
|
} |
|
|
|
jobService.markAsFailed(tenantId, jobId, e.getMessage()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -224,6 +246,25 @@ public class DefaultJobManager implements JobManager { |
|
|
|
Thread.sleep(statsProcessingInterval); // todo: test with bigger interval
|
|
|
|
} |
|
|
|
|
|
|
|
private void sendJobFinishedNotification(Job job) { |
|
|
|
NotificationTemplate template = DefaultNotifications.DefaultNotification.builder() |
|
|
|
.name("Job finished") |
|
|
|
.subject("${type} ${status}") |
|
|
|
.text("${description} ${status}: ${result}") |
|
|
|
.build().toTemplate(); |
|
|
|
GeneralNotificationInfo info = new GeneralNotificationInfo(Map.of( |
|
|
|
"type", job.getType().getTitle(), |
|
|
|
"description", job.getDescription(), |
|
|
|
"status", job.getStatus().name().toLowerCase(), |
|
|
|
"result", job.getResult().getDescription() |
|
|
|
)); |
|
|
|
notificationCenter.sendGeneralWebNotification(job.getTenantId(), new TenantAdministratorsFilter(), template, info); |
|
|
|
} |
|
|
|
|
|
|
|
private JobProcessor getJobProcessor(JobType jobType) { |
|
|
|
return jobProcessors.get(jobType); |
|
|
|
} |
|
|
|
|
|
|
|
@PreDestroy |
|
|
|
private void destroy() { |
|
|
|
jobStatsConsumer.stop(); |
|
|
|
|