Built-in Thread Pools in Java
The java.util.concurrent package provides several factory methods via the Executors class to create common thread pool configurations:
ExecutorService cachedPool = Executors.newCachedThreadPool(); // Scales dynamically with submitted tasks
ExecutorService fixedPool = Executors.newFixedThreadPool(4); // Fixed-size thread pool
ExecutorService singlePool = Executors.newSingleThreadExecutor(); // Single-threaded executor
ExecutorService scheduledPool = Executors.newScheduledThreadPool(4); // Supports delayed or periodic task execution
When submitting tasks in a loop, care must be taken with variable capture. For example, the following code has a bug:
for (int i = 0; i < 10000; i++) {
service.submit(() -> System.out.println(i + Thread.currentThread().getName()));
}
The variable i is not effectively final and changes during loop execution. The correct approach captures a stable copy:
for (int i = 0; i < 10000; i++) {
int taskId = i;
service.submit(() -> System.out.println(taskId + Thread.currentThread().getName()));
}
Custom Thread Pool Implementation
A minimal thread pool requires:
- A collection of worker threads
- A queue to hold pending tasks (
Runnableinstances) - A
submitmethod to enqueue tasks
Here’s a basic implementation:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
class SimpleThreadPool {
private final BlockingQueue<Runnable> taskQueue;
public SimpleThreadPool(int workerCount) {
this.taskQueue = new ArrayBlockingQueue<>(1000);
for (int i = 0; i < workerCount; i++) {
Thread worker = new Thread(() -> {
while (true) {
try {
Runnable task = taskQueue.take();
task.run();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
worker.start();
}
}
public void submit(Runnable task) throws InterruptedException {
taskQueue.put(task);
}
}
An enhanced version supports dynamic thread scaling up to a maximum limit:
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
class ScalableThreadPool {
private final BlockingQueue<Runnable> taskQueue = new ArrayBlockingQueue<>(1000);
private final List<Thread> workers = new ArrayList<>();
private final int maxThreads;
public ScalableThreadPool(int initialThreads, int maxThreads) {
this.maxThreads = maxThreads;
for (int i = 0; i < initialThreads; i++) {
startWorker();
}
}
private void startWorker() {
Thread worker = new Thread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
Runnable task = taskQueue.take();
task.run();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
worker.start();
workers.add(worker);
}
public void submit(Runnable task) throws InterruptedException {
taskQueue.put(task);
if (taskQueue.size() >= 500 && workers.size() < maxThreads) {
startWorker();
}
}
}
Rejection policies can be implemented in submit() by checking queue saturation (e.g., when size exceeds a threshold like 900). Options include:
- Throwing an exception
- Executing the task directly in the calling thread
- Discarding oldest or newest task
Custom Timer Implementation
A timer schedules tasks for future execution. Java’s built-in Timer class demonstrates this:
Timer timer = new Timer();
timer.schedule(new TimerTask() {
public void run() {
System.out.println("Executed after 3 seconds");
}
}, 3000);
To build a custom timer:
- Define a schedulable task with execution time
- Use a priority queue ordered by execution time
- Run a dedicated thread that processes tasks at their scheduled times
Implementation details:
import java.util.PriorityQueue;
class ScheduledTask implements Comparable<ScheduledTask> {
private final Runnable payload;
private final long scheduledTime;
public ScheduledTask(Runnable payload, long delayMs) {
this.payload = payload;
this.scheduledTime = System.currentTimeMillis() + delayMs;
}
public void execute() { payload.run(); }
public long getScheduledTime() { return scheduledTime; }
@Override
public int compareTo(ScheduledTask other) {
return Long.compare(this.scheduledTime, other.scheduledTime);
}
}
class CustomTimer {
private final PriorityQueue<ScheduledTask> queue = new PriorityQueue<>();
private final Object lock = new Object();
public CustomTimer() {
Thread dispatcher = new Thread(() -> {
while (true) {
synchronized (lock) {
if (queue.isEmpty()) {
continue;
}
ScheduledTask next = queue.peek();
long now = System.currentTimeMillis();
if (now >= next.getScheduledTime()) {
queue.poll().execute();
} else {
try {
lock.wait(next.getScheduledTime() - now);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}
}
});
dispatcher.setDaemon(true);
dispatcher.start();
}
public void schedule(Runnable task, long delayMs) {
synchronized (lock) {
queue.offer(new ScheduledTask(task, delayMs));
lock.notify();
}
}
}