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

如何判断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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 03:18:24