为何使用parallelStream分批发送请求时,测试代码能并行执行但业务代码却串行?
为何使用parallelStream分批发送请求时,测试代码能并行执行但业务代码却串行?
嗨,我来帮你拆解下这个问题,这种“测试能跑通并行,业务代码却串行”的情况其实挺常见的,核心原因在于测试代码和业务代码的执行场景、依赖的底层组件完全不同,具体可以从这几个方向排查:
1. 下游服务/客户端的串行限制(最可能的原因)
你的测试代码只是做本地字符串拼接,完全是CPU密集型的本地操作,parallelStream可以轻松分配不同线程并行处理;但业务代码里调用的productPriceService.getProductsPrices(batch)是远程调用(比如请求OpenSearch),如果底层客户端的配置是串行的,那即使你用parallelStream发起调用,实际请求还是会排队执行:
- 比如OpenSearch的RestHighLevelClient如果没配置足够的连接池,默认可能是单连接或者少量连接,导致请求必须等前一个完成才能拿到新连接发送下一个;
- 或者客户端本身用了同步阻塞的单线程模式,比如某些封装的HTTP客户端默认只开一个线程处理请求,自然会串行。
2. parallelStream的线程池适配问题
parallelStream默认使用ForkJoinPool.commonPool(),这个线程池的线程数默认是CPU核心数-1,适合CPU密集型任务,但你的业务场景是IO密集型(远程调用):
- 如果你的机器CPU核心数少(比如2核),那commonPool只有1个工作线程,自然会串行执行;
- 测试代码是CPU密集型,线程数刚好能利用起来,所以能看到并行效果,但IO密集型任务需要更多线程来处理等待时间,commonPool的线程数就不够用了。
3. 业务代码内部的同步/阻塞逻辑
如果productPriceService.getProductsPrices()方法内部有同步锁(比如用了synchronized关键字,或者依赖了单例的同步资源),那即使parallelStream分配了多个线程,也会因为锁的竞争变成串行执行。
排查&解决建议
给你几个实用的排查步骤和优化方案:
- 先确认线程是否真的并行:在业务代码的
map(batch -> ...)里加日志,打印当前线程ID(Thread.currentThread().getId()),如果多个批次的线程ID不同,说明parallelStream确实在并行,问题出在下游客户端;如果线程ID相同,说明parallelStream没并行,可能是线程池的问题。 - 检查客户端连接池配置:比如OpenSearch客户端的
maxConnTotal和maxConnPerRoute参数,把它们调大(比如设成10或20),确保有足够的连接支持并行请求。 - 换用自定义线程池处理IO任务:放弃依赖parallelStream的commonPool,改用
CompletableFuture+自定义线程池,更适合IO密集型场景,示例代码如下:
// 自定义线程池,根据IO任务需求设置线程数 ExecutorService ioExecutor = Executors.newFixedThreadPool(10); List<CompletableFuture<List<ProductPriceResponse>>> futures = productBatches.stream() .map(batch -> CompletableFuture.supplyAsync( () -> productPriceService.getProductsPrices(batch), ioExecutor )) .collect(Collectors.toList()); // 等待所有请求完成,合并结果 List<ProductPriceResponse> productPrices = futures.stream() .map(CompletableFuture::join) .flatMap(List::stream) .collect(Collectors.toList()); // 记得关闭线程池 ioExecutor.shutdown();
- 模拟IO延迟测试:把
productPriceService.getProductsPrices()换成一个带延迟的本地方法(比如Thread.sleep(1000)),如果这时候parallelStream能并行执行,那就100%确定是下游客户端的配置问题。
备注:内容来源于stack exchange,提问作者Yujie Huang
相关产品推荐
相关产品推荐

