NOTE
线程池数目估算
1. 估算线程数目 1.1. 吞吐量=并发度/响应时间 根据吞吐量=并发度/响应时间,那么线程数=QPS 响应时间。 压测 二分法设置线程数,压测查看CPU利用率在80% IO密集型/CPU密集型 如果是计算密集型应用,那么线程数=CPU核心数+1 如果是IO密集型应用,那么线程数=CPU核心数 (1+IO阻塞时间/CPU时间) 1.2. TPS估算 假设要求一个系统的TPS(Transaction Per Second或者Task P
这是历史学习笔记,可能存在过时或不完整的理解。
1. 估算线程数目
1.1. 吞吐量=并发度/响应时间
- 根据吞吐量=并发度/响应时间,那么线程数=QPS*响应时间。
- 压测
- 二分法设置线程数,压测查看CPU利用率在80%
- IO密集型/CPU密集型
- 如果是计算密集型应用,那么线程数=CPU核心数+1
- 如果是IO密集型应用,那么线程数=CPU核心数*(1+IO阻塞时间/CPU时间)
1.2. TPS估算
假设要求一个系统的TPS(Transaction Per Second或者Task Per Second)至少为20,然后假设每个Transaction由一个线程完成,继续假设平均每个线程处理一个Transaction的时间为4s。那么问题转化为:如何设计线程池大小,使得可以在1s内处理完20个Transaction?
计算过程很简单,每个线程的处理能力为0.25TPS,那么要达到20TPS,显然需要20/0.25=80个线程。
需求:每秒处理20个Transaction 条件:一个线程执行Transaciton的时间为4s,那么每个线程的处理能力为1/4 = 0.25TPS(或者说这个接口现实情况是多少s响应) 结果:20/0.25=80
转换成生活问题: 每秒中需要生产20个产品,一个人每秒钟只能生产0.25个产品,需要多少个人才能达到这个结果
转换成总价单价问题: 超市里苹果单价0.25,总价是20,数量有多少个?20/0.25=80
1.3. IO密集 or CPU密集
- IO密集型:IO 操作比 Cpu计算耗时要久的多。线程池大小设置为2N+1
- CPU密集型:Cpu 计算比IO操作耗时要久的多。线程池大小设置为N+1
1.4. Dark Magic
- 抽象类PoolSizeCalculator
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.BlockingQueue;
public abstract class PoolSizeCalculator
{
/**
* The sample queue size to calculate the size of a single {@link Runnable}
* element.
*/
private final int SAMPLE_QUEUE_SIZE = 1000;
/**
* Accuracy of test run. It must finish within 20ms of the testTime
* otherwise we retry the test. This could be configurable.
*/
private final int EPSYLON = 20;
/**
* Control variable for the CPU time investigation.
*/
private volatile boolean expired;
/**
* Time (millis) of the test run in the CPU time calculation.
*/
private final long testtime = 3000;
/**
* Calculates the boundaries of a thread pool for a given {@link Runnable}.
*
* @param targetUtilization the desired utilization of the CPUs (0 <= targetUtilization <= 1)
* @param targetQueueSizeBytes the desired maximum work queue size of the thread pool (bytes)
*/
protected void calculateBoundaries(BigDecimal targetUtilization, BigDecimal targetQueueSizeBytes)
{
calculateOptimalCapacity(targetQueueSizeBytes);
Runnable task = creatTask();
// start(task);
start(task); // warm up phase
long cputime = getCurrentThreadCPUTime();
start(task); // test intervall
cputime = getCurrentThreadCPUTime() - cputime;
long waittime = (testtime * 1000000) - cputime;
calculateOptimalThreadCount(cputime, waittime, targetUtilization);
}
private void calculateOptimalCapacity(BigDecimal targetQueueSizeBytes)
{
long mem = calculateMemoryUsage();
BigDecimal queueCapacity = targetQueueSizeBytes.divide(new BigDecimal(mem), RoundingMode.HALF_UP);
System.out.println("Target queue memory usage (bytes): " + targetQueueSizeBytes);
System.out.println("createTask() produced " + creatTask().getClass().getName() + " which took " + mem + " bytes in a queue");
System.out.println("Formula: " + targetQueueSizeBytes + " / " + mem);
System.out.println("* Recommended queue capacity (bytes): " + queueCapacity);
}
/**
* Brian Goetz' optimal thread count formula, see 'Java Concurrency in * Practice' (chapter 8.2) * * @param cpu * cpu time consumed by considered task * @param wait * wait time of considered task * @param targetUtilization * target utilization of the system
*/
private void calculateOptimalThreadCount(long cpu, long wait, BigDecimal targetUtilization)
{
BigDecimal waitTime = new BigDecimal(wait);
BigDecimal computeTime = new BigDecimal(cpu);
BigDecimal numberOfCPU = new BigDecimal(Runtime.getRuntime().availableProcessors());
BigDecimal optimalthreadcount = numberOfCPU.multiply(targetUtilization).multiply(new BigDecimal(1).add(waitTime.divide(computeTime, RoundingMode.HALF_UP)));
System.out.println("Number of CPU: " + numberOfCPU);
System.out.println("Target utilization: " + targetUtilization);
System.out.println("Elapsed time (nanos): " + (testtime * 1000000));
System.out.println("Compute time (nanos): " + cpu);
System.out.println("Wait time (nanos): " + wait);
System.out.println("Formula: " + numberOfCPU + " * " + targetUtilization + " * (1 + " + waitTime + " / " + computeTime + ")");
System.out.println("* Optimal thread count: " + optimalthreadcount);
}
/**
* Runs the {@link Runnable} over a period defined in {@link #testtime}. * Based on Heinz Kabbutz' ideas * (http://www.javaspecialists.eu/archive/Issue124.html). * * @param task * the runnable under investigation
*/
public void start(Runnable task)
{
long start = 0;
int runs = 0;
do
{
if (++runs > 5)
{
throw new IllegalStateException("Test not accurate");
}
expired = false;
start = System.currentTimeMillis();
Timer timer = new Timer();
timer.schedule(new TimerTask()
{
public void run()
{
expired = true;
}
}, testtime);
while (!expired)
{
task.run();
}
start = System.currentTimeMillis() - start;
timer.cancel();
}
while (Math.abs(start - testtime) > EPSYLON);
collectGarbage(3);
}
private void collectGarbage(int times)
{
for (int i = 0; i < times; i++)
{
System.gc();
try
{
Thread.sleep(10);
}
catch (InterruptedException e)
{
Thread.currentThread().interrupt();
break;
}
}
}
/**
* Calculates the memory usage of a single element in a work queue. Based on
* Heinz Kabbutz' ideas
* (http://www.javaspecialists.eu/archive/Issue029.html).
*
* @return memory usage of a single {@link Runnable} element in the thread
* pools work queue
*/
public long calculateMemoryUsage()
{
BlockingQueue queue = createWorkQueue();
for (int i = 0; i < SAMPLE_QUEUE_SIZE; i++)
{
queue.add(creatTask());
}
long mem0 = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
long mem1 = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
queue = null;
collectGarbage(15);
mem0 = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
queue = createWorkQueue();
for (int i = 0; i < SAMPLE_QUEUE_SIZE; i++)
{
queue.add(creatTask());
}
collectGarbage(15);
mem1 = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
return (mem1 - mem0) / SAMPLE_QUEUE_SIZE;
}
/**
* Create your runnable task here.
*
* @return an instance of your runnable task under investigation
*/
protected abstract Runnable creatTask();
/**
* Return an instance of the queue used in the thread pool.
*
* @return queue instance
*/
protected abstract BlockingQueue createWorkQueue();
/**
* Calculate current cpu time. Various frameworks may be used here,
* depending on the operating system in use. (e.g.
* http://www.hyperic.com/products/sigar). The more accurate the CPU time
* measurement, the more accurate the results for thread count boundaries.
*
* @return current cpu time of current thread
*/
protected abstract long getCurrentThreadCPUTime();
}
- 实现类SimplePoolSizeCaculatorImpl
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.lang.management.ManagementFactory;
import java.math.BigDecimal;
import java.net.HttpURLConnection;
import java.net.URL;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
public class SimplePoolSizeCaculatorImpl extends PoolSizeCalculator
{
@Override
protected Runnable creatTask()
{
return new AsyncIOTask();
}
@Override
protected BlockingQueue createWorkQueue()
{
return new LinkedBlockingQueue(1000);
}
@Override
protected long getCurrentThreadCPUTime()
{
return ManagementFactory.getThreadMXBean().getCurrentThreadCpuTime();
}
public static void main(String[] args)
{
PoolSizeCalculator poolSizeCalculator = new SimplePoolSizeCaculatorImpl();
poolSizeCalculator.calculateBoundaries(new BigDecimal(1.0), new BigDecimal(100000));
}
}
/**
* 自定义的异步IO任务
*
* @author Will
*/
class AsyncIOTask implements Runnable
{
@Override
public void run()
{
HttpURLConnection connection = null;
BufferedReader reader = null;
try
{
String getURL = "https://www.baidu.com/";
URL getUrl = new URL(getURL);
connection = (HttpURLConnection) getUrl.openConnection();
connection.connect();
reader = new BufferedReader(new InputStreamReader(connection.getInputStream()));
String line;
while ((line = reader.readLine()) != null)
{
// empty loop
}
}
catch (IOException e)
{
}
finally
{
if (reader != null)
{
try
{
reader.close();
}
catch (Exception e)
{
}
}
connection.disconnect();
}
}
}
- 输出结果
Target queue memory usage (bytes): 100000 //Queue的总内存占用
createTask() produced com.example.test.AsyncIOTask which took 40 bytes in a queue //queue中的每个task占用40B
Formula: 100000 / 40 //计算Queue大小的公式:通过Queue的总内存占用 / 每个task占用 == queue的大小
* Recommended queue capacity (bytes): 2500 //*****队列的长度为2500*****
Number of CPU: 6 //CPU核数
Target utilization: 1 //CPU目标利用率
Elapsed time (nanos): 3000000000 //总耗时
Compute time (nanos): 62500000 //计算耗时
Wait time (nanos): 2937500000 //等待耗时
Formula: 6 * 1 * (1 + 2937500000 / 62500000) //计算线程数目的公式:CPU核数 * CPU目标利用率 * (CPU目标利用率 + 等待耗时 / 计算耗时)
* Optimal thread count: 288//*****建议的线程池线程数目是288*****
- 构造线程池
ThreadPoolExecutor pool =
new ThreadPoolExecutor(288, 288, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(2500));
1.4.1. 解析
- Queue占用的资源是内存,所以Queue的capacity得通过内存计算出
- 公式:Queue capacity = 希望队列满的时候最大的占用内存 / Queue中每个Element的大小
- 线程占用的资源是CPU,所以线程数得通过CPU消耗率计算出
- 公式:最佳线程数目 = (线程等待时间与线程CPU时间之比 + 1)* CPU数目 * CPU目标利用率
讨论
使用 GitHub 账号参与讨论,评论会保存在 GitHub Issues 中。在 GitHub 查看