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

Java创建于 更新于 约 2 分钟读完historical

这是历史学习笔记,可能存在过时或不完整的理解。

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目标利用率

2. 参考

讨论

使用 GitHub 账号参与讨论,评论会保存在 GitHub Issues 中。在 GitHub 查看