Spring Boot中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

