Код: Выделить всё
'org.springframework.boot' version '3.3.13'
Код: Выделить всё
public class TaskRunner implements DelayedTaskRunner {
private static final String TASK_WORKER_THREAD_NAME_PATTERN = "Task-Worker-%d";
private TaskExecutor taskExecutor;
private ThreadPoolExecutor runnerThreadPool;
public TaskRunner(
TaskExecutor taskExecutor,
Integer defaultThreadPoolSize,
Integer maximumThreadPoolSize,
Duration keepExtraThreadsAliveForSeconds
) {
this.taskExecutor = taskExecutor;
this.runnerThreadPool = new ThreadPoolExecutor(
defaultThreadPoolSize,
maximumThreadPoolSize,
keepExtraThreadsAliveForSeconds.toMillis(),
TimeUnit.MILLISECONDS,
new ArrayBlockingQueue(100),
new ThreadFactoryBuilder().setNameFormat(TASK_WORKER_THREAD_NAME_PATTERN).build()
);
}
public TaskInitiationResult runWorker() {
try {
Future taskRef = runnerThreadPool.submit(this::executeTask);
return new TaskInitiationResult(taskRef);
} catch (RejectedExecutionException e) {
return new TaskInitiationResult(null);
}
}
private void executeTask() {
try {
taskExecutor.executeNextTask();
} catch (Exception e) {
log.error("Error while executing Task", e);
}
}
public void shutdown() {
ExecutorServiceManagementUtil.shutdown(runnerThreadPool, log);
}
@Override
public String getAssignedTaskScope() {
return "TASK_SCOPE_NAME";
}
}
Код: Выделить всё
public class TheTaskExecutor implements TaskExecutor {
private TaskWorkerReportCollector workerReportCollector;
private TaskRepository taskRepository;
private TransactionTemplate transactionTemplate;
private Duration sleepOnWorkNotFound;
private CommandBus commandBus;
private Integer amountOfTasksToSelect;
public TheTaskExecutor(
CommandBus commandBus,
TaskWorkerReportCollector workerReportCollector,
TaskRepository taskRepository,
TransactionTemplate transactionTemplate,
Duration sleepOnWorkNotFound,
Integer amountOfTasksToSelect
) {
this.commandBus = commandBus;
this.amountOfTasksToSelect = amountOfTasksToSelect;
this.taskRepository = taskRepository;
this.transactionTemplate = transactionTemplate;
this.sleepOnWorkNotFound = sleepOnWorkNotFound;
this.workerReportCollector = workerReportCollector;
}
public void executeNextTask() {
long threadId = Thread.currentThread().threadId();
Instant resumeWorkTimestamp = workerReportCollector.getResumeWorkTimestamps().get(threadId);
if (resumeWorkTimestamp != null && resumeWorkTimestamp.isAfter(Instant.now())) {
return;
}
transactionTemplate.executeWithoutResult(transactionStatus -> {
List tasks = taskRepository.deleteWithSelectNextTasksToExecute(amountOfTasksToSelect, Instant.now());
if (tasks.isEmpty()) {
workerReportCollector.put(threadId, Instant.now().plusMillis(sleepOnWorkNotFound.toMillis()));
return;
}
workerReportCollector.evict(threadId);
for (DelayedTask task : tasks) {
try {
commandBus.execute(task.getCommand());
} catch (Exception e) {
throw e;
}
}
});
}
}
Код: Выделить всё
workerReportCollectorПоэтому я ищу общий ответ на вопрос: где поток в полностью изолированном пуле может получить соединение с активной транзакцией? Означает ли это, что Соединение с активной транзакцией остается в пуле Хикари?