Java中如何实现函数并行执行并满足API限制要求?
Java并行任务执行的顺序与并行性保障
我尝试在Java中并行执行两个函数,但不确定当前代码是否满足特定需求。以下是我的代码(ExecutorService相关代码位于主函数中,无关代码已省略):
private static final ExecutorService executorService = Executors.newFixedThreadPool(2); if (!weekly || rowData.get("M/W").equals("W")) { if (callData != null) { Future<?> bidPriceFuture = executorService.submit(() -> { getBidPrice(rowData.get("Symbol"), callData, premiumC, percentMinC, "Call", rowData.get("Call EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Call EPR")), rowData); }); Future<?> updateVerticalsFuture = executorService.submit(() -> { updateVerticals(rowData.get("Symbol"), callData); }); bidPriceFuture.get(); updateVerticalsFuture.get(); } if (putData != null) { System.out.println("Running the Put Stuff"); Future<?> putFuture = executorService.submit(() -> { getBidPrice(rowData.get("Symbol"), putData, premiumP, percentMinP, "Put", rowData.get("Put EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Put EPR")), rowData); }); Future<?> updateVerticalPut = executorService.submit(() -> { updateVerticals(rowData.get("Symbol"), putData); }); putFuture.get(); updateVerticalPut.get(); }
需求总结
- 第一个if分支(callData非空)中的
getBidPrice和updateVerticals需要并行执行 - 第二个if分支(putData非空)的所有任务必须等待第一个分支的所有任务完成后再启动(受API限制)
- 当前通过
Future.get()控制顺序,但担心该操作会破坏并行性,希望得到真正并行的实现方案
解决方案
你当前的代码其实没有破坏callData分支内的并行性——两个任务都已提交到线程池后才调用get(),它们会在线程池的两个线程里并行执行,get()只是阻塞当前主线程等待任务完成。不过可以优化写法,让逻辑更清晰,同时严格保障put分支的执行顺序:
优化方案1:批量提交+统一等待
把call分支的两个任务先全部提交,一次性等待所有完成后再处理put分支:
private static final ExecutorService executorService = Executors.newFixedThreadPool(2); if (!weekly || rowData.get("M/W").equals("W")) { List<Future<?>> callFutures = new ArrayList<>(); if (callData != null) { // 提交call分支的两个并行任务 callFutures.add(executorService.submit(() -> { getBidPrice(rowData.get("Symbol"), callData, premiumC, percentMinC, "Call", rowData.get("Call EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Call EPR")), rowData); })); callFutures.add(executorService.submit(() -> { updateVerticals(rowData.get("Symbol"), callData); })); } // 等待call分支所有任务完成,再处理put分支 for (Future<?> future : callFutures) { try { future.get(); } catch (InterruptedException | ExecutionException e) { // 实际业务中可替换为日志记录或异常抛出 e.printStackTrace(); } } if (putData != null) { System.out.println("Running the Put Stuff"); List<Future<?>> putFutures = new ArrayList<>(); // 提交put分支的两个并行任务 putFutures.add(executorService.submit(() -> { getBidPrice(rowData.get("Symbol"), putData, premiumP, percentMinP, "Put", rowData.get("Put EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Put EPR")), rowData); })); putFutures.add(executorService.submit(() -> { updateVerticals(rowData.get("Symbol"), putData); })); // 等待put分支任务完成 for (Future<?> future : putFutures) { try { future.get(); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); } } } }
优化方案2:使用CompletableFuture(Java 8+)
如果使用Java 8及以上版本,CompletableFuture能让异步逻辑更简洁,天然支持并行和顺序控制:
private static final ExecutorService executorService = Executors.newFixedThreadPool(2); if (!weekly || rowData.get("M/W").equals("W")) { // 定义call阶段的两个并行任务 CompletableFuture<Void> callPhase = CompletableFuture.allOf( CompletableFuture.runAsync(() -> { if (callData != null) { getBidPrice(rowData.get("Symbol"), callData, premiumC, percentMinC, "Call", rowData.get("Call EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Call EPR")), rowData); } }, executorService), CompletableFuture.runAsync(() -> { if (callData != null) { updateVerticals(rowData.get("Symbol"), callData); } }, executorService) ); // 等待call阶段全部完成后,执行put阶段 callPhase.thenRun(() -> { if (putData != null) { System.out.println("Running the Put Stuff"); CompletableFuture.allOf( CompletableFuture.runAsync(() -> { getBidPrice(rowData.get("Symbol"), putData, premiumP, percentMinP, "Put", rowData.get("Put EPR").equals("#N/A") ? 0 : Double.parseDouble(rowData.get("Put EPR")), rowData); }, executorService), CompletableFuture.runAsync(() -> { updateVerticals(rowData.get("Symbol"), putData); }, executorService) ).join(); // 等待put阶段所有任务完成 } }).join(); }
关键说明
- 并行性保障:只要在调用等待方法(
get()/join())前提交任务,两个任务就会在线程池的独立线程中并行执行,等待操作仅阻塞主线程,不影响任务本身的并行运行。 - 顺序控制:通过先等待call分支所有任务完成,再提交put分支任务,严格满足API限制的顺序要求。
- 异常处理:实际开发中需合理处理
InterruptedException(线程被中断)和ExecutionException(任务执行抛出异常),避免程序静默失败。
内容的提问来源于stack exchange,提问作者Rishab
相关产品推荐
相关产品推荐

