Executor Thread Pool - limit queue size and dequeue oldest

Viewed 10503

I am using a fixed thread pool for a consumer of produced messages within a spring boot application. My producer is producing (a lot) faster than the producer is able to handle a message, therefore the queue of the thread pool seems to be "flooding".

What would be the best way to limit the queue size? The intended queue behaviour would be "if the queue is full, remove the head and insert the new Runnable". Is it possible to configure Executors thread pool like this?

3 Answers

ThreadPoolExecutor supports this function via ThreadPoolExecutor.DiscardOldestPolicy:

A handler for rejected tasks that discards the oldest unhandled request and then retries execute, unless the executor is shut down, in which case the task is discarded.

You need to construct the pool with this policy mannully, for exmaple:

int poolSize = ...;
int queueSize = ...;
RejectedExecutionHandler handler = new ThreadPoolExecutor.DiscardOldestPolicy();

ExecutorService executorService = new ThreadPoolExecutor(poolSize, poolSize,
    0L, TimeUnit.MILLISECONDS,
    new LinkedBlockingQueue<>(queueSize),
    handler);

This will create a thread pool for you of the size that you pass.

ExecutorService service = Executors.newFixedThreadPool(THREAD_SIZE);

This internally creates an instance of ThreadPoolExecutor, which implements ExecutorService.

public static ExecutorService newFixedThreadPool(int nThreads) {
    return new ThreadPoolExecutor(nThreads, nThreads,
                                  0L, TimeUnit.MILLISECONDS,
                                  new LinkedBlockingQueue<Runnable>());
}

To create a custom thead pool, you can just do.

ExecutorService service =   new ThreadPoolExecutor(5, 5,
                                  0L, TimeUnit.MILLISECONDS,
                                  new LinkedBlockingQueue<Runnable>(10));

Here we can specify the size of the queue, using the overloaded constructor of the LinkedBlockingQueue.

    public LinkedBlockingQueue(int capacity) {
        if (capacity <= 0) throw new IllegalArgumentException();
        this.capacity = capacity;
        last = head = new Node<E>(null);
    }

Hope this helps. Cheers !!!

For example if you working with Data Base(psql) which has capable of 100 connections at a time . and task may takes 2000ms ...

        int THREADS = 50;
        ExecutorService exe = new ThreadPoolExecutor(THREADS,
                                                    50,
                                                    0L,
                                                     TimeUnit.MILLISECONDS,
                                                     new ArrayBlockingQueue<>(10),
                                                     new ThreadPoolExecutor.CallerRunsPolicy()); ```
     
Related