You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java中如何高效并发调用层级RESTful API?非阻塞IO可行吗?

高效并发收集部门ID的非阻塞IO方案

针对你的层级部门API遍历需求,普通多线程因OS线程的上下文切换和资源占用开销过高,推荐使用异步非阻塞HTTP请求 + CompletableFuture + 虚拟线程的方案,既能实现高并发,又能大幅降低性能开销。

核心思路

  1. 用异步HTTP客户端发送请求,线程无需等待IO响应,可复用处理其他任务
  2. 借助CompletableFuture实现异步递归遍历部门层级,并行处理子部门请求
  3. 用线程安全集合去重,避免重复请求同一部门
  4. 利用虚拟线程(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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 22:15:36