更换为自定义线程池

v1.4.1
Parker 4 years ago
parent c298e6daf0
commit 7b73dc2093

@ -0,0 +1,53 @@
package org.opsli.common.thread;
import cn.hutool.core.util.StrUtil;
import lombok.extern.slf4j.Slf4j;
/**
* @Author:
* @CreateTime: 2020-10-08 10:24
* @Description: 线
*/
@Slf4j
public class AsyncProcessQueue {
// ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
/**
* Task <br>
* Executor <br>
*/
public static class TaskWrapper implements Runnable {
private final Runnable gift;
public TaskWrapper(final Runnable target) {
this.gift = target;
}
@Override
public void run() {
// 捕获异常,避免在 Executor 里面被吞掉了
if (gift != null) {
try {
gift.run();
} catch (Exception e) {
String errMsg = StrUtil.format("线程池-包装的目标执行异常: {}", e.getMessage());
log.error(errMsg, e);
}
}
}
}
// ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
/**
*
*
* @param task
* @return
*/
public static boolean execute(final Runnable task) {
return AsyncProcessor.executeTask(new TaskWrapper(task));
}
}

@ -0,0 +1,128 @@
package org.opsli.common.thread;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.concurrent.BasicThreadFactory;
import java.util.concurrent.*;
/**
* @Author:
* @CreateTime: 2020-10-08 10:24
* @Description: 线
*/
@Slf4j
public class AsyncProcessor {
/**
* <br>
*/
private static final int DEFAULT_MAX_CONCURRENT = Runtime.getRuntime().availableProcessors() * 2;
/**
* 线
*/
private static final String THREAD_POOL_NAME = "ExternalConvertProcessPool-%d";
/**
* 线
*/
private static final ThreadFactory FACTORY = new BasicThreadFactory.Builder()
.namingPattern(THREAD_POOL_NAME)
.daemon(true).build();
/**
*
*/
private static final int DEFAULT_SIZE = 500;
/**
* 线
*/
private static final int DEFAULT_WAIT_TIME = 10;
/**
* 线
*/
private static final long DEFAULT_KEEP_ALIVE = 60L;
/**NewEntryServiceImpl.java:689
* Executor
*/
private static final ExecutorService EXECUTOR;
/**
*
*/
private static final BlockingQueue<Runnable> EXECUTOR_QUEUE = new ArrayBlockingQueue<>(DEFAULT_SIZE);
static {
// 创建 Executor
// 此处默认最大值改为处理器数量的 4 倍
try {
EXECUTOR = new ThreadPoolExecutor(DEFAULT_MAX_CONCURRENT, DEFAULT_MAX_CONCURRENT * 4, DEFAULT_KEEP_ALIVE,
TimeUnit.SECONDS, EXECUTOR_QUEUE, FACTORY);
// 关闭事件的挂钩
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("AsyncProcessor shutting down.");
EXECUTOR.shutdown();
try {
// 等待1秒执行关闭
if (!EXECUTOR.awaitTermination(DEFAULT_WAIT_TIME, TimeUnit.SECONDS)) {
log.error("AsyncProcessor shutdown immediately due to wait timeout.");
EXECUTOR.shutdownNow();
}
} catch (InterruptedException e) {
log.error("AsyncProcessor shutdown interrupted.");
EXECUTOR.shutdownNow();
}
log.info("AsyncProcessor shutdown complete.");
}));
} catch (Exception e) {
log.error("AsyncProcessor init error.", e);
throw new ExceptionInInitializerError(e);
}
}
/**
*
*/
private AsyncProcessor() {
}
/**
* <br>
* {@link }
*
* @param task
* @return
*/
public static boolean executeTask(Runnable task) {
try {
EXECUTOR.execute(task);
} catch (RejectedExecutionException e) {
log.error("Task executing was rejected.", e);
return false;
}
return true;
}
/**
* <br>
* {@link }
*
* @param task
* @return
*/
public static <T> Future<T> submitTask(Callable<T> task) {
try {
return EXECUTOR.submit(task);
} catch (RejectedExecutionException e) {
log.error("Task executing was rejected.", e);
throw new UnsupportedOperationException("Unable to submit the task, rejected.", e);
}
}
}

@ -13,7 +13,7 @@
* License for the specific language governing permissions and limitations under
* the License.
*/
package org.opsli.api.thread.factory;
package org.opsli.common.thread.factory;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.atomic.AtomicInteger;

@ -2,16 +2,12 @@ package org.opsli.core.thread;
import lombok.extern.slf4j.Slf4j;
import org.opsli.api.base.result.ResultVo;
import org.opsli.api.thread.factory.NameableThreadFactory;
import org.opsli.api.web.system.logs.LogsApi;
import org.opsli.api.wrapper.system.logs.LogsModel;
import org.opsli.common.api.TokenThreadLocal;
import org.opsli.common.thread.AsyncProcessQueue;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/**
* @BelongsProject: tank-design
* @BelongsPackage: com.parker.tank.net.thread
@ -24,10 +20,6 @@ import java.util.concurrent.Executors;
public class LogsThreadPool {
private static final ExecutorService EXECUTOR_SERVICE = Executors.newFixedThreadPool(4,
new NameableThreadFactory("日志保存线程"));
/** 日志API */
private static LogsApi logsApi;
@ -39,16 +31,13 @@ public class LogsThreadPool {
if(logsModel == null){
return;
}
EXECUTOR_SERVICE.submit(()->{
AsyncProcessQueue.execute(()->{
// 存储临时 token
try {
ResultVo<?> ret = logsApi.insert(logsModel);
if(!ret.isSuccess()){
log.error(ret.getMsg());
}
}catch (Exception e){
log.error(e.getMessage(), e);
}
ResultVo<?> ret = logsApi.insert(logsModel);
if(!ret.isSuccess()){
log.error(ret.getMsg());
}
});
}

Loading…
Cancel
Save