订单系统里,一个下单成功事件要同时触发:发短信、更新统计、同步搜索索引、通知风控。一开始我用 @TransactionalEventListener 在一个方法里全干了,结果短信接口抖一下,整个下单事务被拖慢。这种"一个事件多个消费者"的场景,正是 Spring Integration 的主场。
消息通道把生产消费解耦
Spring Integration 的核心是管道(pipe-and-filters):消息从 MessageChannel 进,经过过滤器、路由器、转换器,到终点。下单后我只往一个通道发消息,谁关心谁订阅:
@Autowired
private MessageChannel orderCompletedChannel;
public void onOrder(Order o) {
orderCompletedChannel.send(MessageBuilder.withPayload(o)
.setHeader("orderId", o.getId()).build());
}
这样下单主流程只管"发出去",后面怎么处理它不关心,也不背锅。
过滤器与路由:消息去哪自己说
不是所有消息都要走全部下游。比如只有金额大于 1000 的订单才进风控,用 Filter:
@Filter(inputChannel = "orderCompletedChannel", outputChannel = "riskChannel")
public boolean needRisk(Order o) {
return o.getAmount().compareTo(BigDecimal.valueOf(1000)) > 0;
}
不过滤的消息默认被丢弃,若想"不满足条件也去别处",用 Router 按条件分发到不同通道。我们做过一个路由:VIP 订单走"加急通道"(同步处理),普通订单走"普通通道"(异步削峰),两边是不同的线程池:
@Router(inputChannel = "orderCompletedChannel")
public String route(Order o) {
return o.isVip() ? "urgentChannel" : "normalChannel";
}
与 MQ 集成:把内存管道接到外部
纯内存通道进程一挂消息就没。我们用 Jms 或 Kafka 适配器把通道接到 MQ,下单消息持久化到 Kafka,下游各自消费:
@Bean
public KafkaProducerMessageChannelAdapter kafkaOut(
KafkaTemplate template) {
return new KafkaProducerMessageChannelAdapter<>(template, "order-completed");
}
这里有个关键认知:Spring Integration 的通道是"进程内管道",MQ 适配器是"进程间桥梁"。同一进程内用通道解耦、路由、过滤;跨进程、要持久化和重放的,交给 MQ。两者配合,既灵活又可靠。
一个异步配置的坑
默认通道是同步的,发消息的线程直接执行消费者。一旦消费者慢(发短信 800ms),下单线程被卡。我们给通道加 ExecutorChannel 异步化:
@Bean
public MessageChannel orderCompletedChannel(TaskExecutor exec) {
return new ExecutorChannel(exec);
}
但异步后事务边界变了——消息发出去那一刻下单事务可能还没提交,消费者读到的是旧数据。解决是让消费者在消息里带"事件版本"或延后读,等事务真正提交。我们用了 @TransactionalEventListener(phase = AFTER_COMMIT) 作为发送触发点,保证发消息时事务已落地。
小结
Spring Integration 把"事件怎么分发、谁来处理、同步还是异步"用声明式通道表达清楚,比手写一堆 @Async 方法可读、可测。过滤器做取舍、路由器做分发、MQ 适配器做跨进程持久化,三件套覆盖了大多数事件驱动需求。提醒一句:进程内通道不持久化,关键事件一定要接到 MQ 上,否则重启即丢。它不是 MQ 的替代品,而是 MQ 之上的一层编排。