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.