Java中如何高效并发调用层级RESTful API?非阻塞IO可行吗?
高效并发收集部门ID的非阻塞IO方案
针对你的层级部门API遍历需求,普通多线程因OS线程的上下文切换和资源占用开销过高,推荐使用异步非阻塞HTTP请求 + CompletableFuture + 虚拟线程的方案,既能实现高并发,又能大幅降低性能开销。
核心思路
- 用异步HTTP客户端发送请求,线程无需等待IO响应,可复用处理其他任务
- 借助CompletableFuture实现异步递归遍历部门层级,并行处理子部门请求
- 用线程安全集合去重,避免重复请求同一部门
- 利用虚拟线程(Java 19+)进一步降低线程资源占用(低版本Java可替换为优化后的线程池)
具体实现
1. 定义数据模型
先把API返回的JSON映射为Java类:
import com.fasterxml.jackson.databind.ObjectMapper; public class Division { private int division; private String subdivisions; private boolean status; // Getters & Setters public int getDivision() { return division; } public void setDivision(int division) { this.division = division; } public String getSubdivisions() { return subdivisions; } public void setSubdivisions(String subdivisions) { this.subdivisions = subdivisions; } public boolean isStatus() { return status; } public void setStatus(boolean status) { this.status = status; } }
2. 异步遍历实现
import com.fasterxml.jackson.core.JsonProcessingException; import java.net.URI; import java.util.Arrays; import java.util.List; import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; import java.util.stream.Collectors; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; public class DepartmentTreeFetcher { // 异步HTTP客户端,用虚拟线程池执行器(Java 19+) private static final HttpClient ASYNC_HTTP_CLIENT = HttpClient.newBuilder() .executor(Executors.newVirtualThreadPerTaskExecutor()) .build(); private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); // 线程安全集合:记录已处理的部门ID,避免重复请求 private final ConcurrentHashMap<Integer, Boolean> processedDivisions = new ConcurrentHashMap<>(); // 可选:API限流信号量,控制并发请求数,避免触发API限流 private final Semaphore requestSemaphore = new Semaphore(15); public CompletableFuture<Void> fetchAllDepartments(int rootDeptId) { return processDepartment(rootDeptId); } private CompletableFuture<Void> processDepartment(int deptId) { // 先判断是否已处理过该部门,避免重复请求 if (processedDivisions.putIfAbsent(deptId, Boolean.TRUE) != null) { return CompletableFuture.completedFuture(null); } // 先获取限流许可 return CompletableFuture.runAsync(() -> { try { requestSemaphore.acquire(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Failed to acquire request permit", e); } }) .thenCompose(ignored -> { // 构建API请求 URI requestUri = URI.create("http://your-api-domain/getDivision/dID=" + deptId); HttpRequest request = HttpRequest.newBuilder(requestUri) .GET() .build(); // 异步发送请求并处理响应 return ASYNC_HTTP_CLIENT.sendAsync(request, HttpResponse.BodyHandlers.ofString()) .whenComplete((response, error) -> requestSemaphore.release()) // 无论成败都释放许可 .thenApply(HttpResponse::body) .thenApply(this::parseDivisionResponse) .thenCompose(this::processSubDepartments); }); } private Division parseDivisionResponse(String responseBody) { try { return OBJECT_MAPPER.readValue(responseBody, Division.class); } catch (JsonProcessingException e) { throw new RuntimeException("Failed to parse division API response", e); } } private CompletableFuture<Void> processSubDepartments(Division division) { // 拆分子部门ID列表 List<Integer> subDeptIds = Arrays.stream(division.getSubdivisions().split(",")) .map(String::trim) .filter(idStr -> !idStr.isBlank()) .map(Integer::parseInt) .collect(Collectors.toList()); // 并行异步处理所有子部门 List<CompletableFuture<Void>> subTasks = subDeptIds.stream() .map(this::processDepartment) .collect(Collectors.toList()); // 等待所有子部门处理完成 return CompletableFuture.allOf(subTasks.toArray(new CompletableFuture[0])); } // 获取收集到的所有部门ID public Set<Integer> getAllDepartmentIds() { return processedDivisions.keySet(); } // 测试入口 public static void main(String[] args) throws Exception { DepartmentTreeFetcher fetcher = new DepartmentTreeFetcher(); // 从顶层部门ID=1开始遍历 fetcher.fetchAllDepartments(1).get(); System.out.println("All collected department IDs: " + fetcher.getAllDepartmentIds()); } }
方案优势
- 非阻塞IO:异步HTTP请求不会阻塞线程,线程可复用处理其他请求,大幅减少线程上下文切换开销
- 虚拟线程:Java 19+的虚拟线程由JVM调度,无需占用OS线程资源,能支撑上万级别的并发请求,内存占用极低
- 自动去重:用
ConcurrentHashMap的putIfAbsent确保每个部门只被请求一次,避免无效请求 - 可限流:通过
Semaphore控制并发请求数,适配API的QPS限制,避免触发限流或报错
低版本Java兼容调整
如果使用Java 8-18,无法使用虚拟线程,可以将HTTP客户端的执行器替换为优化后的线程池:
private static final HttpClient ASYNC_HTTP_CLIENT = HttpClient.newBuilder() .executor(Executors.newFixedThreadPool(10)) // 根据API性能调整线程数 .build();
内容的提问来源于stack exchange,提问作者Andrew Law
相关产品推荐
相关产品推荐

