Java CompletionService多线程:如何将响应结果映射到对应请求
问题解决方法
你当前的代码无法关联请求和响应,核心原因是提交任务后没有保存任务和原请求的绑定关系,Future完全可以解决这个问题,有两种常用实现方案,按需选择即可。
方案1:通过Future做映射绑定(改动最小)
CompletionService.submit()方法本身就会返回和任务绑定的Future对象,你只需要在提交任务时,把返回的Future和对应的原请求存在映射表里,后续从take()拿到完成的Future时,直接查表就能拿到对应的原请求。
参考实现代码:
ExecutorService myThreadPoolExecutor = Executors.newFixedThreadPool(3); // 注意修正原代码的泛型错误,原代码写的FeeEstimationResult和实际返回的MyResponse不匹配 CompletionService<MyResponse> completionService = new ExecutorCompletionService<>(myThreadPoolExecutor); // 存储Future和对应原请求的映射 Map<Future<MyResponse>, MyRequest> futureRequestMap = new HashMap<>(); for (final MyRequest myRequest : myRequestList) { Future<MyResponse> taskFuture = completionService.submit(() -> getMyResponse(myRequest)); futureRequestMap.put(taskFuture, myRequest); } for (int i = 0; i < myRequestList.size(); i++) { try { Future<MyResponse> completedFuture = completionService.take(); // 取出对应请求后直接移除映射,避免内存泄漏 MyRequest matchedRequest = futureRequestMap.remove(completedFuture); MyResponse myResponse = completedFuture.get(); // 后续可直接使用matchedRequest和对应myResponse做业务处理 } catch (final Exception e) { log.error("批量请求处理异常", e); } } // 线程池使用完记得shutdown,避免线程泄漏 myThreadPoolExecutor.shutdown();
方案2:包装返回结果(无需维护映射表,更省心)
如果可以调整任务的返回值结构,完全不用额外维护Map,直接让任务返回「请求+响应」的包装对象,拿到结果的时候自然就能对应上原请求,出错概率更低。
首先定义一个简单的包装类:
// 可根据自己的项目习惯替换成lombok注解或者手写构造、get/set public class RequestResponseHolder { private MyRequest request; private MyResponse response; public RequestResponseHolder(MyRequest request, MyResponse response) { this.request = request; this.response = response; } public MyRequest getRequest() {return request;} public MyResponse getResponse() {return response;} }
调整任务提交和结果获取逻辑:
ExecutorService myThreadPoolExecutor = Executors.newFixedThreadPool(3); CompletionService<RequestResponseHolder> completionService = new ExecutorCompletionService<>(myThreadPoolExecutor); for (final MyRequest myRequest : myRequestList) { completionService.submit(() -> { MyResponse response = getMyResponse(myRequest); return new RequestResponseHolder(myRequest, response); }); } for (int i = 0; i < myRequestList.size(); i++) { try { RequestResponseHolder resultHolder = completionService.take().get(); MyRequest matchedRequest = resultHolder.getRequest(); MyResponse myResponse = resultHolder.getResponse(); // 直接做后续业务处理即可 } catch (final Exception e) { log.error("批量请求处理异常", e); } } myThreadPoolExecutor.shutdown();
选型参考
- 不想改动原有
getMyResponse方法的返回逻辑,选方案1,改动量最小,注意处理完结果及时移除Map里的条目即可 - 可以调整返回结构的场景优先选方案2,不需要额外维护映射关系,不会出现漏存、漏删映射导致的逻辑bug
另外提醒原代码的两个小问题:
CompletionService声明的泛型是FeeEstimationResult,实际取结果用的是MyResponse,泛型不统一会有类型转换异常风险- 线程池使用完没有调用shutdown方法,会导致核心线程一直存活,可能引发内存泄漏、应用无法正常退出的问题
内容的提问来源于stack exchange,提问作者zxwang
相关产品推荐
相关产品推荐

