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

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();
}

关键说明

  1. 并行性保障:只要在调用等待方法(get()/join())前提交任务,两个任务就会在线程池的独立线程中并行执行,等待操作仅阻塞主线程,不影响任务本身的并行运行。
  2. 顺序控制:通过先等待call分支所有任务完成,再提交put分支任务,严格满足API限制的顺序要求。
  3. 异常处理:实际开发中需合理处理InterruptedException(线程被中断)和ExecutionException(任务执行抛出异常),避免程序静默失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:00:19