TaskPoolConfig.java 2.79 KB
package com.infoloop.tianting.config;

import com.infoloop.tianting.context.LoginContextHolder;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;
import org.springframework.core.task.TaskDecorator;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;

import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;

import static com.infoloop.tianting.constant.ConfigConstants.ASYNC_EXECUTOR;
import static com.infoloop.tianting.constant.ConfigConstants.SCHEDULED_TASK_EXECUTOR;

@Slf4j
@Configuration
@EnableScheduling
@EnableAsync
public class TaskPoolConfig {

    @Primary
    @Bean(SCHEDULED_TASK_EXECUTOR)
    public TaskScheduler scheduledExecutorService() {
        final var scheduler = new ThreadPoolTaskScheduler();
        scheduler.setPoolSize(10);
        scheduler.setThreadNamePrefix("Scheduled-Task-Executor-");
        scheduler.setWaitForTasksToCompleteOnShutdown(true);
        scheduler.setAwaitTerminationSeconds(60);
        scheduler.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        scheduler.setErrorHandler(e -> {
            // 记录异常日志或发送通知
            log.error("Scheduled task error", e);
        });
        scheduler.initialize();
        return scheduler;
    }

    @Bean(ASYNC_EXECUTOR)
    public Executor taskExecutor() {
        final var executor = new ThreadPoolTaskExecutor();
        final var availableProcessors = Runtime.getRuntime().availableProcessors();
        executor.setCorePoolSize(availableProcessors * 2 + 1);
        executor.setMaxPoolSize(50);
        executor.setQueueCapacity(200);
        executor.setKeepAliveSeconds(60);
        executor.setThreadNamePrefix("Async-Executor-");
        executor.setTaskDecorator(new ThreadLocalTaskDecorator());
        executor.initialize();
        return executor;
    }

    @SuppressWarnings("all")
    public class ThreadLocalTaskDecorator implements TaskDecorator {
        @Override
        public Runnable decorate(Runnable runnable) {
            // 保存当前线程的登录信息
            final var loginOperatorInfo = LoginContextHolder.getLoginInfo();
            return () -> {
                try {
                    LoginContextHolder.setLoginInfo(loginOperatorInfo);
                    runnable.run();
                } finally {
                    LoginContextHolder.clear();
                }
            };
        }
    }

}