如何判断Spring Integration Flow运行状态、避免重复触发及模拟竞态测试?
解决Spring Integration流程并发触发问题
一、是否有内置方式判断流程是否运行?
Spring Integration 没有直接提供线程安全的方法来判断当前是否有流程实例在执行。StandardIntegrationFlow.isRunning() 仅表示该流是否处于启动状态,而非是否有任务正在处理,因此无法满足你的需求。你需要自行实现全局的执行锁机制。
二、正确使用AtomicBoolean控制整个流程并发
核心思路是:用AtomicBoolean维护一个全局运行标志,在流程入口(Cron轮询/HTTP触发)原子性获取锁,在整个批次的所有消息处理完成后释放锁,避免单条消息处理完就释放导致的并发问题。
1. 实现全局锁服务
@Component public class FlowExecutionGuard { private final AtomicBoolean isFlowRunning = new AtomicBoolean(false); // 原子性尝试获取锁,成功返回true,失败返回false public boolean tryAcquire() { return isFlowRunning.compareAndSet(false, true); } // 释放锁 public void release() { isFlowRunning.set(false); } }
2. 在HTTP触发入口添加锁判断
@RestController public class FlowTriggerController { private final FlowExecutionGuard guard; private final MessageChannel triggerChannel; public FlowTriggerController(FlowExecutionGuard guard, MessageChannel triggerChannel) { this.guard = guard; this.triggerChannel = triggerChannel; } @PostMapping("/trigger-flow") public ResponseEntity<String> triggerFlow() { if (!guard.tryAcquire()) { return ResponseEntity.status(HttpStatus.CONFLICT).body("流程正在运行中"); } try { // 发送触发消息,携带锁实例到消息头 Message<?> triggerMsg = MessageBuilder.withPayload("manual-trigger") .setHeader("flowGuard", guard) .build(); triggerChannel.send(triggerMsg); return ResponseEntity.ok("流程已启动"); } catch (Exception e) { // 异常时释放锁 guard.release(); throw new RuntimeException("触发流程失败", e); } } }
3. 在Cron轮询入口添加锁判断
@Bean @InboundChannelAdapter(value = "directChannel", poller = @Poller(cron = "0 */5 * * * ?")) public MessageSource<List<Data>> dbPoller(FlowExecutionGuard guard, JdbcTemplate jdbcTemplate) { return () -> { if (!guard.tryAcquire()) { // 流程正在运行,跳过本次轮询 return null; } try { List<Data> dataList = jdbcTemplate.query("SELECT * FROM target_table", new DataRowMapper()); return MessageBuilder.withPayload(dataList) .setHeader("flowGuard", guard) .build(); } catch (Exception e) { guard.release(); throw new RuntimeException("轮询数据失败", e); } }; }
4. 聚合拆分消息,完成后释放锁
因为你的流程包含消息拆分,必须等所有拆分后的消息处理完成再释放锁,否则会出现提前解锁导致的并发问题。通过Aggregator组件实现批次聚合:
@Bean public IntegrationFlow mainFlow() { return IntegrationFlows.from("directChannel") // 拆分消息 .split() // 转换数据 .transform(Data.class, this::transformToTarget) // 提交到线程池处理 .channel(MessageChannels.executor(Executors.newFixedThreadPool(5))) // 执行业务操作 .handle((payload, headers) -> { // 你的业务逻辑 return payload; }) // 聚合所有拆分后的消息,确保批次处理完成 .aggregate(aggregatorSpec -> aggregatorSpec .correlationStrategy(msg -> msg.getHeaders().getId()) .releaseStrategy(group -> group.getMessages().size() == group.getSequenceSize())) // 释放锁(添加异常处理确保锁必释放) .handle((payload, headers) -> { FlowExecutionGuard guard = headers.get("flowGuard", FlowExecutionGuard.class); guard.release(); return payload; }, spec -> spec.advice(new GuardReleaseAdvice())) // 输出到目标通道 .channel("outputChannel") .get(); } // 统一处理锁释放的Advice,确保异常时也能解锁 private static class GuardReleaseAdvice extends AbstractRequestHandlerAdvice { @Override protected Object doInvoke(ExecutionCallback callback, Object target, Message<?> message) throws Exception { try { return callback.execute(); } finally { FlowExecutionGuard guard = message.getHeaders().get("flowGuard", FlowExecutionGuard.class); if (guard != null) { guard.release(); } } } }
三、模拟竞态条件测试
通过多线程并发触发,结合CountDownLatch和@SpyBean验证锁机制的有效性:
@SpringBootTest @AutoConfigureMockMvc public class FlowConcurrencyTest { @Autowired private MockMvc mockMvc; @SpyBean private FlowExecutionGuard flowGuard; @Test void testConcurrentTriggerPrevented() throws InterruptedException { int concurrentCount = 5; CountDownLatch latch = new CountDownLatch(concurrentCount); ExecutorService executor = Executors.newFixedThreadPool(concurrentCount); // 提交并发请求 for (int i = 0; i < concurrentCount; i++) { executor.submit(() -> { try { mockMvc.perform(post("/trigger-flow")); } catch (Exception e) { e.printStackTrace(); } finally { latch.countDown(); } }); } // 等待所有请求完成 latch.await(10, TimeUnit.SECONDS); // 验证锁仅被成功获取1次,释放1次 verify(flowGuard, times(concurrentCount)).tryAcquire(); verify(flowGuard, times(1)).release(); } }
你也可以通过验证业务操作的执行次数(比如数据库写入次数)来确认流程仅执行一次。
内容的提问来源于stack exchange,提问作者Sal
相关产品推荐
相关产品推荐

