Implementing Dynamic Thread Pools for Parallel Multi-Provider API Queries

Scenario Overview

Consider a requirement where data from multiple telecom providers (such as China Telecom, China Unicom, and China Mobile) needs to be aggregated and displayed on a single page. This involves calling three distinct API endpoints and consolidating the results.

Serial Query Implementation

A straightforward approach involves executing requests sequentially:

JSONObject telecomResult = httpClient.get(telecomEndpoint, parameters, headers);
JSONObject unicomResult = httpClient.get(unicomEndpoint, parameters, headers);
JSONObject mobileResult = httpClient.get(mobileEndpoint, parameters, headers);

While sequential execution works when requests have dependencies, it introduces several performance bottlenecks when queries are independent:

Increased Latency: Each request must complete before the next begins. Total response time equals the sum of all individual response times, resulting in noticeable delays especially under poor network conditions.

Limited Throughput: The system processes only one request at a time, restricting overall throughput. In high-concurrency environments, this underutilizes available processing capacity.

Resource Underutilization: CPU and other resources remain idle while waiting for I/O operations to complete. Modern multi-core systems cannot leverage parallel processing capabilities.

Poor Scalability: Execution time grows linearly with the number of queries. Adding hardware resources provides minimal performance improvement, creating a scalability ceiling.

Parallel Execution Strategy

To overcome these limitations, parallel query execution should be employed. Using Java's CompletableFuture API enables concurrent HTTP requests with a single blocking call to wait for all results:

CompletableFuture<JSONObject> telecomFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(telecomEndpoint, parameters, headers);
});

CompletableFuture<JSONObject> unicomFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(unicomEndpoint, parameters, headers);
});

CompletableFuture<JSONObject> mobileFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(mobileEndpoint, parameters, headers);
});

CompletableFuture.allOf(telecomFuture, unicomFuture, mobileFuture).join();

However, this approach relies on the default common fork-join pool, which lacks configurability and monitoring capabilities. For production systems, implemmenting a custom thread pool provides better control.

Dynamic Thread Pool with Nacos Integration

A configurable thread pool based on Nacos allows runtime adjustment of pool parameters and real-time status monitoring. The following implementation demonstrates dynamic thread pool management:

import com.alibaba.cloud.nacos.NacosConfigManager;
import com.alibaba.cloud.nacos.NacosConfigProperties;
import com.alibaba.nacos.api.config.listener.Listener;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.context.annotation.Configuration;

import java.util.concurrent.Executor;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionHandler;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

@RefreshScope
@Configuration
public class ConfigurableThreadPool implements InitializingBean {
    
    @Value("${thread.pool.core-size}")
    private String corePoolSize;
    
    @Value("${thread.pool.maximum-size}")
    private String maximumPoolSize;
    
    private static ThreadPoolExecutor executorService;
    
    @Autowired
    private NacosConfigManager configManager;
    
    @Autowired
    private NacosConfigProperties configProperties;
    
    @Override
    public void afterPropertiesSet() {
        executorService = new ThreadPoolExecutor(
            Integer.parseInt(corePoolSize),
            Integer.parseInt(maximumPoolSize),
            10L,
            TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(10),
            new ThreadFactoryBuilder().setNameFormat("async-worker-%d").build(),
            new RejectedExecutionHandler() {
                @Override
                public void rejectedExecution(Runnable task, ThreadPoolExecutor pool) {
                    System.err.println("Task rejected: queue full or shutdown");
                }
            }
        );
        
        configManager.getConfigService().addListener(
            "application-prod.yml",
            configProperties.getGroup(),
            new Listener() {
                @Override
                public Executor getExecutor() {
                    return null;
                }
                
                @Override
                public void receiveConfigInfo(String configContent) {
                    updatePoolConfiguration(Integer.parseInt(corePoolSize), 
                                           Integer.parseInt(maximumPoolSize));
                }
            }
        );
    }
    
    public String getPoolStatus() {
        return String.format(
            "Core: %d, Active: %d, Max: %d, Queue: %d, TotalTasks: %d",
            executorService.getCorePoolSize(),
            executorService.getActiveCount(),
            executorService.getMaximumPoolSize(),
            executorService.getQueue().size(),
            executorService.getTaskCount()
        );
    }
    
    public void submitTask(Runnable task) {
        executorService.execute(task);
    }
    
    private void updatePoolConfiguration(int coreSize, int maxSize) {
        executorService.setCorePoolSize(coreSize);
        executorService.setMaximumPoolSize(maxSize);
    }
}

Applying this thread pool to parallel API queries:

CompletableFuture<JSONObject> telecomFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(telecomEndpoint, parameters, headers);
}, ConfigurableThreadPool.executorService);

CompletableFuture<JSONObject> unicomFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(unicomEndpoint, parameters, headers);
}, ConfigurableThreadPool.executorService);

CompletableFuture<JSONObject> mobileFuture = CompletableFuture.supplyAsync(() -> {
    return httpClient.get(mobileEndpoint, parameters, headers);
}, ConfigurableThreadPool.executorService);

CompletableFuture.allOf(telecomFuture, unicomFuture, mobileFuture).join();

Thread Pool Sizing Guidelines

Optimal pool size depends on task characteristics:

I/O-Bound Tasks: These tasks spend significant time waiting for external resources (network, database). Larger thread counts enable better resource utilization since threads can proceed when others are blocked. Setting pool size to 2 * CPU cores or higher often yields optimal throughput for I/O-heavy workloads.

CPU-Bound Tasks: Computation-intensive tasks consume CPU cycles continuously. Excessive threads cause contention for CPU resources, degrading performance. Pool size of CPU cores + 1 or CPU cores + 2 balances utilization without excessive context switching overhead.

Tags: thread-pool java parallel-processing CompletableFuture Nacos

Posted on Tue, 06 Oct 2026 16:26:41 +0000 by Instigate