Fundamentals of Concurrency
Program, Process, and Thread
- Program: A static collection of instructions written to perform a specific task. It is passive code residing on disk.
- Process: An active execution of a program. It represents a running instance with a lifecycle (start, run, terminate).
- Thread: A lightweight, independent path of execution within a process. A process can contain multiple threads.
Understanding Thread Pools
Context: Constantly creating and destroying threads is resource-intensive and harms performance, especially under high concurrency.
Solution: Pre-initialize a "pool" of threads. When a task arrives, it borrows a thread from the pool; when finished, the thread returns to the pool for reuse. This is analogous to a taxi stand where cars are ready and waiting for passengers.
Advantages:
- Speed: Eliimnates the latency of thread creation.
- Efficiency: Reuses existing threads, lowering memory overhead.
- Management: Offers control over system resources via specific parameters:
corePoolSize: The base number of threads kept alive.maximumPoolSize: The max threads allowed when the queue is full.keepAliveTime: Timeout for idle non-core threads before termination.
Spring Boot Async Implementation
While standard Java threads can be created by extending the Thread class or implementing Runnable, Spring Boot simplifies this via the @Async annotation and a managed task executor.
Configuring the Executor Bean
Define a configuration class to customize the thread pool behavior. This allows you to define core size, queue capacity, and naming conventions.
@Configuration
@EnableAsync
public class TaskPoolConfig {
@Bean(name = "customTaskExecutor")
public Executor createTaskExecutor() {
ThreadPoolTaskExecutor poolExecutor = new ThreadPoolTaskExecutor();
// Base number of threads
poolExecutor.setCorePoolSize(5);
// Max threads when queue is full
poolExecutor.setMaxPoolSize(15);
// Queue size before creating new threads beyond core
poolExecutor.setQueueCapacity(200);
// Idle timeout for threads exceeding core size
poolExecutor.setKeepAliveSeconds(30);
// Prefix for thread names for easier debugging
poolExecutor.setThreadNamePrefix("async-worker-");
// Policy when queue is full (e.g., discard or caller runs)
poolExecutor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
poolExecutor.initialize();
return poolExecutor;
}
}
Defining Asynchronous Methods
Annotate service methods with @Async to execute them in a separate thread managed by the configured executor.
@Slf4j
@Service
public class BackgroundService {
// Explicitly linking to the bean defined above
@Async("customTaskExecutor")
public void processData(String payload) {
log.info("Processing payload: {} in thread: {}", payload, Thread.currentThread().getName());
try {
TimeUnit.MILLISECONDS.sleep(500);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("Task interrupted", e);
}
}
}
Critical Usage Constraints
- Bean Requirement: The class containing
@Asyncmethods must be a Spring-managed bean (e.g.,@Service,@Component). - Enable Annotation: Ensure
@EnableAsyncis present in your configuration. - No Static Methods:
@Asyncdoes not work onstaticmethods. - Self-Invocation Issue: Calling an
@Asyncmethod from within the same class will not work because the call bypasses the Spring proxy. The method must be called from a different bean. - Return Types: Asynchronous methods should return
void,Future<T>, orCompletableFuture<T>. Returning other types results in the method running async but returningnull.
Handling Results and Exceptions
Retrieving Results with CompletableFuture
Java 8's CompletableFuture provides a robust way to handle async results.
@Async("customTaskExecutor")
public CompletableFuture<String> fetchDataAsync(int id) {
log.info("Fetching data for id: {}", id);
// Simulate work
Thread.sleep(1000);
return CompletableFuture.completedFuture("Data-" + id);
}
Caller side:
CompletableFuture<String> task1 = backgroundService.fetchDataAsync(1);
CompletableFuture<String> task2 = backgroundService.fetchDataAsync(2);
// Combine results or wait for all
CompletableFuture<Void> allTasks = CompletableFuture.allOf(task1, task2);
allTasks.get(); // Block until all complete
log.info("Result 1: {}", task1.get());
Exception Handling Strategies
-
Internal Try-Catch: The simplest way is to handle exceptions inside the asynchronous method itself.
-
Via Future/CompletableFuture: Catch exceptiosn when calling
.get(). ```try { String result = task1.get(); } catch (ExecutionException e) { log.error("Error inside async task: {}", e.getCause().getMessage()); } -
CompletableFuture Exceptionally: Use functional style to provide fallbacks. ```
CompletableFuture<String> safeTask = backgroundService.fetchDataAsync(99) .exceptionally(throwable -> { log.error("Task failed, returning fallback.", throwable); return "Fallback-Data"; });