Java并发递归爬取好友列表的最优实现方案咨询
Java并发递归解析好友列表的性能瓶颈分析与优化方案
现有代码的性能瓶颈分析
第一个版本(递归并行)的问题
- 线程池资源不足:默认
CompletableFuture.supplyAsync()复用ForkJoinPool.commonPool(),该线程池线程数等于CPU核心数,而网络请求是IO密集型操作,线程长期阻塞等待响应,导致大量任务排队,无法充分利用并发。 - 异常处理不严谨:将
IOException包装为IllegalStateException抛出,线程池中的未捕获异常可能被吞掉,且无日志记录,难以排查问题。 - 无去重机制:好友列表存在重复用户时,会重复发起网络请求,浪费带宽和服务器资源。
- 集合类型强制转换不安全:
Collectors.toList()返回的列表不一定是ArrayList,强制转换可能抛出ClassCastException。
第二个版本(DFS串行分层)的问题
- 分层串行执行:必须等待当前层所有用户的好友列表全部获取完成,才会进入下一层递归,无法利用"当前层部分任务完成后立即启动下一层请求"的并发潜力,整体执行时间等于各层执行时间之和。
- 异常处理随意:直接调用
printStackTrace(),异常信息分散且容易丢失,返回null会导致后续flatMap操作抛出空指针异常。 - 同第一个版本的线程池、去重、集合转换问题。
优化后的并发实现方案
核心优化点
- 使用自定义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; } }
优化效果说明
- 并发利用率提升:自定义线程池适配IO密集场景,避免线程阻塞导致的任务排队;并行BFS模型让下一层请求提前启动,减少整体等待时间。
- 资源浪费减少:去重机制避免重复请求同一用户,降低带宽和服务器压力。
- 稳定性增强:严谨的异常处理和线程中断处理,避免异常丢失或程序挂死。
内容的提问来源于stack exchange,提问作者user12420288
相关产品推荐
相关产品推荐

