NOTE

Estimating Thread Pool Size

1. Estimate the number of threads 1.1. Throughput = concurrency / response time 1.2. TPS estimation 1.3. I/O-intensive or CPU-intensive 1.4. Dark Magic 2. References

JavaCreated Updated 2 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. Estimate the Number of Threads

1.1. Throughput = Concurrency / Response Time

  • Based on throughput = concurrency / response time, the number of threads = QPS * response time.
  • Load testing
    • Use a binary-search approach to set the thread count and run load tests, checking that CPU utilization is around 80%.
  • I/O-intensive / CPU-intensive
    • For compute-intensive applications, number of threads = number of CPU cores + 1.
    • For I/O-intensive applications, number of threads = number of CPU cores * (1 + I/O blocking time / CPU time).

1.2. TPS Estimation

Suppose a system is required to achieve at least 20 TPS (Transactions Per Second or Tasks Per Second). Assume each transaction is completed by one thread, and that on average each thread takes 4 seconds to process one transaction. The problem then becomes: how should the thread-pool size be designed so that 20 transactions can be processed within one second?

The calculation is simple. Each thread has a processing capacity of 0.25 TPS, so to reach 20 TPS, 20 / 0.25 = 80 threads are required.

Requirement: process 20 transactions per second.

Condition: one thread takes 4 seconds to execute one transaction, so each thread has a processing capacity of 1 / 4 = 0.25 TPS (or, in a real system, use the actual response time of the interface).

Result: 20 / 0.25 = 80.

Convert it into an everyday-life problem:

20 products need to be produced each second. One person can produce only 0.25 products per second. How many people are needed to achieve the target?

Convert it into a total-price/unit-price problem:

Apples in a supermarket cost 0.25 each and the total price is 20. How many apples are there? 20 / 0.25 = 80.

1.3. I/O-Intensive or CPU-Intensive

  • I/O-intensive: I/O operations take much longer than CPU computation. Set the thread-pool size to 2N + 1.
  • CPU-intensive: CPU computation takes much longer than I/O operations. Set the thread-pool size to N + 1.

1.4. Dark Magic

  • Abstract class 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();

}
  • Implementation class 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));
    }

}

/**
 * Custom asynchronous I/O task.
 *
 * @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();
        }

    }

}
  • Output
Target queue memory usage (bytes): 100000 // Total memory occupied by the queue
createTask() produced com.example.test.AsyncIOTask which took 40 bytes in a queue // Each task in the queue occupies 40 B
Formula: 100000 / 40 // Queue-size formula: total queue memory / memory occupied by each task = queue size
* Recommended queue capacity (bytes): 2500 // *****The queue length is 2500*****
Number of CPU: 6 // Number of CPU cores
Target utilization: 1 // Target CPU utilization
Elapsed time (nanos): 3000000000 // Total elapsed time
Compute time (nanos): 62500000 // Compute time
Wait time (nanos): 2937500000 // Wait time
Formula: 6 * 1 * (1 + 2937500000 / 62500000) // Thread-count formula: CPU cores * target CPU utilization * (1 + wait time / compute time)
* Optimal thread count: 288// *****The recommended number of threads in the pool is 288*****
  • Construct the thread pool
ThreadPoolExecutor pool =
 new ThreadPoolExecutor(288, 288, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(2500));

1.4.1. Analysis

  • The resource occupied by the Queue is memory, so the Queue capacity must be calculated from memory.
    • Formula: Queue capacity = maximum memory you want the full queue to occupy / size of each Element in the Queue.
  • Threads consume CPU resources, so the number of threads must be calculated from CPU utilization.
    • Formula: optimal number of threads = (thread wait time / thread CPU time + 1) * number of CPUs * target CPU utilization.

2. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub