线程池工具类
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
// 线程池构建器 模板用法参见 common.ThreadPoolBuilderTest
// 队列满了之后才会创建第(corePoolSize+1)个线程, 而LinkedBlockingQueue 默认大小为int.Max,SynchronousQueue 大小为1
// 默认队列满之后的拒绝策略是抛出异常, 会导致任务线程停止并且pool.shutdown()永远不能正常结束
// 必须捕获任务线程的异常
// 线程按顺序消费队列
public final class ThreadPoolUtil {
private static final Logger log = LoggerFactory.getLogger(ThreadPoolUtil.class);
// cpu核心数
public static final int cpu_num = Runtime.getRuntime().availableProcessors();
public static ThreadPoolExecutor buildPool(String name) {
return buildPool(name, cpu_num);
}
/**
* 构建线程池
*
* @param threadNamePrefix 任务线程名字前缀
* @param maxSize 线程池最大线程数
* @return 线程池
*/
public static ThreadPoolExecutor buildPool(String threadNamePrefix, int maxSize) {
// 创建固定线程数,任务队列无穷大的线程池
ThreadPoolExecutor pool = new ThreadPoolExecutor(maxSize, maxSize, 6, TimeUnit.MINUTES, new LinkedBlockingQueue());
// 允许核心线程被回收,使线程池空闲时会收缩至0
pool.allowCoreThreadTimeOut(true);
// 设置任务线程的名字
pool.setThreadFactory(new SimpleThreadFactory(threadNamePrefix));
// 设置任务拒绝策略(若不配置任务队列大小, 其实也就不存在拒绝的情况)
pool.setRejectedExecutionHandler(new SimpleRejectedExecutionHandler());
return pool;
}
// 默认拒绝策略 ThreadPoolExecutor.AbortPolicy
static class SimpleRejectedExecutionHandler implements RejectedExecutionHandler {
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
log.error("Task {} rejected from {}", r.toString(), e.toString());
}
}
// 参考Executors.defaultThreadFactory()
static class SimpleThreadFactory implements ThreadFactory {
private final ThreadGroup group;
private final AtomicInteger threadNumber = new AtomicInteger(1);
private final String namePrefix;
SimpleThreadFactory(String namePrefix) {
SecurityManager s = System.getSecurityManager();
group = (s != null) ? s.getThreadGroup() : Thread.currentThread().getThreadGroup();
this.namePrefix = namePrefix + "-";
}
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(group, r, namePrefix + threadNumber.getAndIncrement(), 0);
if (t.isDaemon()) t.setDaemon(false);
if (t.getPriority() != Thread.NORM_PRIORITY) t.setPriority(Thread.NORM_PRIORITY);
return t;
}
}
}
线程池使用模板代码
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
// 线程池使用模板代码
public class ThreadPoolUtilTest {
private static final Logger log = LoggerFactory.getLogger(ThreadPoolUtilTest.class);
private static final String poolName = "test";
private static final ThreadPoolExecutor pool = ThreadPoolUtil.buildPool(poolName, 4);
// 不阻塞调用线程
@Test
public void noBlockingExecute() {
List
cook
- 线程池业务逻辑必须使用try-catch包裹, 否则异常发生时日志中不会有任何异常信息, 不利于异常时排查
- ThreadPoolExecutor#getCompletedTaskCount 方法使用的ReentrantLock锁, 读和写都会加锁, 测试发现分别使用 每次任务完成后使用getCompletedTaskCount 计算进度 和 新建线程定时计算进度ThreadPoolUtil#progressMonitor 两种方式时, 整个批次完成耗时相差无几
服务器托管,北京服务器托管,服务器租用 http://www.fwqtg.net
机房租用,北京机房租用,IDC机房托管, http://www.fwqtg.net
相关推荐: Doris写入数据异常提示actual column number in csv file is less than schema column number
版本信息: Flink 1.17.1 Doris 1.2.3 Flink Doris Connector 1.4.0 写入方式 采用 String 数据流,依照社区网站的样例代码,在sink之前将数据转换为DataStream,分隔符采用”t”。 运行异常 通…