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

Spring Boot中CompletableFuture未等待前置函数执行完成问题排查

Spring Boot中Modbus半双工通信的CompletableFuture同步问题排查

问题背景

  • 现有startPooling函数每10秒调用requestForState获取设备当前状态,sendSignal函数用于发送信号修改设备参数(如温度);
  • 因Modbus是半双工通信,sendSignal需等待startPooling所在线程执行完毕,但实际sendSignal并未等待requestForState的CompletableFuture执行完成;
  • 测试示例可正常运行,但实际业务场景中问题复现,需排查原因。

相关代码实现

requestForState函数代码

if (this.updatePending.get()) {
    throw new ModbusBusyException("Cannot request for state device " + getId() + " due to previous request has not end yet");
}
this.setRunningTask(
        this.getRunningTask()
                .thenRun(() -> this.updatePending.set(true))
                .<request-for-state-then-compose-to-state-class>
                .whenComplete((state, error) -> this.updatePending.set(false)))
return this.getRunningTask().join();

sendSignal函数代码

if(modbusDevice.getUpdatePending().get()){
    LOG.warn("Resync {} is still working ",modbusDevice.getDevice().getId());
    // 此处join未生效
    modbusDevice.getRunningTask().join();
}
modbusDevice.setRunningTask(modbusDevice.getRunningTask()
        .thenRun(() -> modbusDevice.getUpdatePending().set(true))
        .thenCompose(voidSignal -> <Comple>)
        .thenRun(() -> modbusDevice.getUpdatePending().set(false))
        .thenRun(this.getState(modbusDevice))).join();

测试示例代码

public static void main(String[] args) throws InterruptedException {
    for (int i = 0; i < 10; i++) {
        new Thread(() -> System.out.println("STATE"+ modbusDevice.requestForState())).start();

        Thread.sleep(50);

        new Thread(() -> {
            System.out.println(LocalDateTime.now() + " - IS THE END OF RESYNC ? " + modbusDevice.getRunningTask().isDone());
            modbusDevice.setRunningTask(
                    modbusDevice.getRunningTask()
                            .thenRun(() -> modbusDevice.getUpdatePending().set(true))
                            .thenRun(() -> System.out.println(Thread.currentThread().getId() + ":\t" + LocalDateTime.now() + " - START SENDING EXECUTE"))
                            .thenCompose(voidSignal -> modbusDevice.setFanSpeed(JohnsonControlsFanSpeed.MEDIUM))
                            .thenRun(() -> modbusDevice.getUpdatePending().set(false))
                            .thenRun(() -> System.out.println(Thread.currentThread().getId() + ":\t" + LocalDateTime.now() + " - STOP SENDING EXECUTE"))
            ).join();
        }).start();

        Thread.sleep(5000);
    }
}

可能的排查方向

1. 竞态条件导致状态判断失效

sendSignal中先判断updatePending.get()再调用runningTask.join()的操作不是原子性的。如果在判断和join之间,requestForState刚好完成并重置updatePending,同时有新任务提交,会导致sendSignal跳过等待直接提交新任务,破坏半双工的顺序执行逻辑。

2. runningTask引用的非原子更新问题

在requestForState和sendSignal中,基于旧的runningTask引用构建新任务链后调用setRunningTask,如果多线程同时获取旧引用并更新,会出现后提交的任务覆盖前一个的情况,丢失任务依赖关系,导致无法等待前置任务完成。

建议改用AtomicReference的updateAndGet方法实现原子更新:

this.runningTask.updateAndGet(task -> task
        .thenRun(() -> this.updatePending.set(true))
        .thenCompose(...)
        .whenComplete((state, error) -> this.updatePending.set(false)));

3. 线程池执行策略与join()阻塞冲突

测试示例使用手动创建的线程,而Spring Boot中CompletableFuture默认使用ForkJoinPool.commonPool()。如果任务中频繁调用join()阻塞线程,可能引发线程池饥饿,导致任务调度顺序偏离预期,出现半双工通信的同步问题。

建议为Modbus通信单独创建专用线程池,避免与其他任务共用线程资源,同时尽量避免在链式任务中使用join()阻塞,改用thenRunAsync指定专用线程池执行任务。

4. updatePending状态修改时机不严谨

requestForState中thenRun(() -> this.updatePending.set(true))是在旧任务完成后执行,若旧任务已完成则会立即执行,此时若有其他线程正在判断updatePending,会导致状态判断错误。另外,whenComplete中重置状态时,若任务抛出异常,可能出现状态重置不及时或错误的情况。

建议将updatePending.set(true)绑定到实际Modbus请求的前置操作,而非thenRun中;同时用handle方法统一处理任务成功/失败后的状态重置逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 04:53:09