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

Java并发递归爬取好友列表的最优实现方案咨询

Java并发递归解析好友列表的性能瓶颈分析与优化方案

现有代码的性能瓶颈分析

第一个版本(递归并行)的问题

  1. 线程池资源不足:默认CompletableFuture.supplyAsync()复用ForkJoinPool.commonPool(),该线程池线程数等于CPU核心数,而网络请求是IO密集型操作,线程长期阻塞等待响应,导致大量任务排队,无法充分利用并发。
  2. 异常处理不严谨:将IOException包装为IllegalStateException抛出,线程池中的未捕获异常可能被吞掉,且无日志记录,难以排查问题。
  3. 无去重机制:好友列表存在重复用户时,会重复发起网络请求,浪费带宽和服务器资源。
  4. 集合类型强制转换不安全:Collectors.toList()返回的列表不一定是ArrayList,强制转换可能抛出ClassCastException。

第二个版本(DFS串行分层)的问题

  1. 分层串行执行:必须等待当前层所有用户的好友列表全部获取完成,才会进入下一层递归,无法利用"当前层部分任务完成后立即启动下一层请求"的并发潜力,整体执行时间等于各层执行时间之和。
  2. 异常处理随意:直接调用printStackTrace(),异常信息分散且容易丢失,返回null会导致后续flatMap操作抛出空指针异常。
  3. 同第一个版本的线程池、去重、集合转换问题。

优化后的并发实现方案

核心优化点

  • 使用自定义IO密集型线程池:线程数设为CPU核心数的3-4倍,适配网络IO阻塞场景。
  • 加入全局去重机制:用ConcurrentHashMap.newKeySet()记录已处理的用户URI,避免重复请求。
  • 采用并行BFS模型:每层任务完成一个就立即启动下一层请求,而非等待整层完成,最大化并发利用率。
  • 严谨的异常处理:用日志记录异常,避免异常丢失或吞掉。
  • 避免不安全的集合操作:指定明确的集合类型,减少不必要的复制。

代码实现

1. 初始化自定义线程池

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class FriendFetcher {
    private static final Logger log = LoggerFactory.getLogger(FriendFetcher.class);
    // 自定义IO密集型线程池:线程数=CPU核心数*4,后台线程避免阻止JVM退出
    private static final ExecutorService IO_THREAD_POOL = Executors.newFixedThreadPool(
            Runtime.getRuntime().availableProcessors() * 4,
            new ThreadFactory() {
                private final AtomicInteger counter = new AtomicInteger(0);
                @Override
                public Thread newThread(Runnable r) {
                    Thread thread = new Thread(r);
                    thread.setName("friend-fetcher-io-" + counter.incrementAndGet());
                    thread.setDaemon(true);
                    return thread;
                }
            }
    );
    // 线程安全的去重集合
    private final Set<String> processedUris = ConcurrentHashMap.newKeySet();

2. 并行BFS实现(推荐)

public List<String> fetchAllFriends(String startUri, int maxDepth) throws InterruptedException, ExecutionException {
        List<String> allFriends = Collections.synchronizedList(new ArrayList<>());
        Queue<CompletableFuture<List<String>>> currentLevelTasks = new LinkedList<>();

        // 初始化第一层任务
        if (processedUris.add(startUri)) {
            currentLevelTasks.add(CompletableFuture.supplyAsync(() -> {
                try {
                    List<String> friends = getUsers(startUri);
                    allFriends.addAll(friends);
                    friends.forEach(processedUris::add);
                    return friends;
                } catch (IOException e) {
                    log.error("获取用户[{}]的好友列表失败", startUri, e);
                    return Collections.emptyList();
                }
            }, IO_THREAD_POOL));
        }

        int currentDepth = 1;
        while (currentDepth < maxDepth && !currentLevelTasks.isEmpty()) {
            Queue<CompletableFuture<List<String>>> nextLevelTasks = new LinkedList<>();

            // 等待当前层所有任务完成,同时处理结果并生成下一层任务
            CompletableFuture<Void> currentLevelDone = CompletableFuture.allOf(
                    currentLevelTasks.toArray(new CompletableFuture[0])
            ).thenRunAsync(() -> {
                currentLevelTasks.forEach(task -> {
                    try {
                        List<String> currentFriends = task.get();
                        for (String friendUri : currentFriends) {
                            if (processedUris.add(friendUri)) { // 仅处理未访问过的用户
                                CompletableFuture<List<String>> nextTask = CompletableFuture.supplyAsync(() -> {
                                    try {
                                        List<String> nextFriends = getUsers(friendUri);
                                        allFriends.addAll(nextFriends);
                                        nextFriends.forEach(processedUris::add);
                                        return nextFriends;
                                    } catch (IOException e) {
                                        log.error("获取用户[{}]的好友列表失败", friendUri, e);
                                        return Collections.emptyList();
                                    }
                                }, IO_THREAD_POOL);
                                nextLevelTasks.add(nextTask);
                            }
                        }
                    } catch (InterruptedException | ExecutionException e) {
                        log.error("处理任务结果失败", e);
                        Thread.currentThread().interrupt();
                    }
                });
            }, IO_THREAD_POOL);

            currentLevelDone.get();
            currentLevelTasks = nextLevelTasks;
            currentDepth++;
        }

        // 若为一次性任务,执行完成后关闭线程池;长期运行则保留
        // IO_THREAD_POOL.shutdown();
        return allFriends;
    }

    // 原有的网络请求方法
    private List<String> getUsers(String uri) throws IOException {
        // 模拟长期网络调用
        // return 从网站获取的好友URI列表
    }
}

3. 优化后的递归并行实现

public List<String> fetchFriendsRecursive(String startUri, int maxDepth) throws IOException {
        processedUris.clear(); // 重置去重集合
        return doFetchFriends(startUri, maxDepth, 1);
    }

    private List<String> doFetchFriends(String uri, int maxDepth, int currentDepth) {
        if (currentDepth > maxDepth || !processedUris.add(uri)) {
            return Collections.emptyList();
        }

        List<String> currentFriends;
        try {
            currentFriends = getUsers(uri);
        } catch (IOException e) {
            log.error("获取用户[{}]的好友列表失败", uri, e);
            return Collections.emptyList();
        }

        List<String> allFriends = new ArrayList<>(currentFriends);
        processedUris.addAll(currentFriends);

        if (currentDepth == maxDepth) {
            return allFriends;
        }

        // 并行处理所有好友的递归请求
        List<CompletableFuture<List<String>>> futures = currentFriends.stream()
                .map(friendUri -> CompletableFuture.supplyAsync(
                        () -> doFetchFriends(friendUri, maxDepth, currentDepth + 1),
                        IO_THREAD_POOL
                ))
                .collect(Collectors.toList());

        // 等待所有子任务完成并合并结果
        try {
            CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).get();
            futures.stream()
                    .map(task -> {
                        try {
                            return task.get();
                        } catch (InterruptedException | ExecutionException e) {
                            log.error("获取子任务结果失败", e);
                            Thread.currentThread().interrupt();
                            return Collections.emptyList();
                        }
                    })
                    .flatMap(List::stream)
                    .forEach(allFriends::add);
        } catch (InterruptedException | ExecutionException e) {
            log.error("等待所有子任务完成失败", e);
            Thread.currentThread().interrupt();
        }

        return allFriends;
    }
}

优化效果说明

  1. 并发利用率提升:自定义线程池适配IO密集场景,避免线程阻塞导致的任务排队;并行BFS模型让下一层请求提前启动,减少整体等待时间。
  2. 资源浪费减少:去重机制避免重复请求同一用户,降低带宽和服务器压力。
  3. 稳定性增强:严谨的异常处理和线程中断处理,避免异常丢失或程序挂死。

内容的提问来源于stack exchange,提问作者user12420288

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 08:54:18