ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Spring Event本地好用远程别硬撑:从事件总线到消息队列的可靠演进

Spring Event本地好用远程别硬撑:从事件总线到消息队列的可靠演进 带过分布式的人都体会过一种割裂单体时代用 Spring Event 做业务解耦一行publishEvent、一个EventListener干净利落。等系统一拆有人开始琢磨“能不能让 Spring Event 跨服务跑”于是把本地事件塞进 Redis、塞进 MQ再让远程服务重新发布一次。最初几天确实“好用”代码也能跑在线上的大流量和故障面前问题一个个冒出来。最后你会看到同一个景象本地代码里 Spring Event 照常用但跨服务的“远程 Spring Event”几乎没人再碰大家宁可引入 Kafka、RocketMQ也不愿再维护那套桥接代码。这篇文章就围绕“本地”和“远程”两条线拆一拆讲清楚它为什么在本地好用、搬到远程以后为什么变味以及如果真的要从本地走向远程应该怎么设计才不踩坑。1. Spring Event 在本地为什么好用1.1 一个容器级事件总线的底层逻辑Spring Event 本质上是 Spring 容器提供的一个进程内观察者模式实现。核心成员就三个ApplicationEvent事件对象、ApplicationEventPublisher发布器、ApplicationListener或EventListener监听器。当业务代码调用publisher.publishEvent(...)时容器会在内部找到所有匹配的监听器然后逐个调用。关键在于“进程内”这三个字。因为发布端和监听端在同一个 JVM 里事件对象直接以引用方式传递不需要序列化和反序列化没有网络开销也没有中间件参与。这种设计带来的直接好处有三个一是性能极高一次发布就是一次普通方法调用二是类型安全监听器拿到的就是业务里创建的那个对象字段、泛型都能完整保留三是调试方便异常堆栈会把发布和监听的调用链串在一起IDE 里直接打断点整个链路清清楚楚。我自己在单体项目里最常见的写法就是这样Component public class OrderService { private final ApplicationEventPublisher publisher; public OrderService(ApplicationEventPublisher publisher) { this.publisher publisher; } public void createOrder(OrderDTO dto) { // 业务逻辑校验、落库 OrderCreatedEvent event new OrderCreatedEvent(orderId, userId, amount); publisher.publishEvent(event); } } Component public class CouponListener { EventListener public void onOrderCreated(OrderCreatedEvent event) { // 发放优惠券、发送通知、更新统计 } }这段代码背后Spring 会自动创建ApplicationEventMulticaster默认使用同步执行器在调用线程里完成监听器的回调。也就是说发布事件之后监听器里的逻辑必须执行完才能返回。很多人第一次接触时会觉得“这算什么异步”但恰恰是这个同步特性让事务边界变得可控也让单元测试变得极其简单。你不需要启动消息中间件只需要ApplicationEventPublisher和监听器。1.2 本地事件的三个高级姿势事务、异步、顺序如果只会用最基础的EventListener那本地事件的潜力其实没发挥完。我在实际项目里用得最多的三个扩展能力分别解决三个不同的问题。第一个是TransactionalEventListener。它可以让事件在事务提交后再被处理避免出现“事件先发出去了事务后面回滚了”的脏数据问题。用法是在监听方法上标注事务阶段最常用的是AFTER_COMMITTransactionalEventListener(phase TransactionPhase.AFTER_COMMIT) public void onOrderCreated(OrderCreatedEvent event) { // 此时订单事务已提交可以放心发通知 }这里要特别注意这个语义只在“监听器所在进程内”有效。如果事件经过远程传输发到另一个服务的 JVM 里TransactionalEventListener管的是哪边的事务答案是哪边都不管。很多远程事件踩坑就是从“事务事件监听器失效”开始的。第二个是Async异步化。Spring Event 默认同步执行如果监听器里有发送短信、调用外部接口这类耗时操作会阻塞主流程。搭配Async注解事件就可以丢到线程池里执行Component public class NotificationListener { Async EventListener public void onOrderCreated(OrderCreatedEvent event) { // 异步发送短信不阻塞下单接口 } }使用异步监听器的前提是配置了EnableAsync并且要明白异步后异常不会回到发布者链路追踪也需要手动传播。还有个容易踩的坑同一个类内部调用带有Async的方法会失效因为代理不生效。如果事件监听器内部又调用了同类方法记得把异步逻辑独立到另一个 Bean 里。第三个是Order顺序控制。本地同步事件依赖同一个线程所以监听器之间的相对顺序是确定的。想控制执行顺序直接声明Order(1)、Order(2)EventListener Order(10) public void doFirst(OrderCreatedEvent event) { } EventListener Order(20) public void doSecond(OrderCreatedEvent event) { }这些能力组合在一起在单个应用内部解耦业务时非常顺手。但请注意一个事实这些功能都有一个隐含前提——所有东西都在同一个 JVM 里。一旦事件要跨服务这个前提就没了。2. 把 Spring Event 搬到远程看似简单实则处处是坑2.1 远程 Spring Event 常见的桥接写法Spring 官方其实并没有提供一个开箱即用的“远程事件”功能。大家说的远程 Spring Event通常是自己做桥接在服务 A 里监听本地事件然后把事件内容转发到 Redis、Kafka 或 RabbitMQ服务 B 收到消息后反序列化再通过ApplicationEventPublisher.publishEvent(...)发到本地 Spring 容器从而触发监听器。这种桥接方式刚写出来时非常诱人因为它让业务代码保持了统一我依然用EventListener写监听逻辑感觉只是多了一层转发而已。一段极简的实现思路大致长这样Component public class RemoteEventForwarder { private final RedisTemplateString, Object redisTemplate; private final ObjectMapper objectMapper; EventListener public void onLocalEvent(OrderCreatedEvent event) throws JsonProcessingException { String payload objectMapper.writeValueAsString(event); // 发送到 Redis Stream redisTemplate.opsForStream().add( StreamRecords.newRecord() .ofObject(payload) .withStreamKey(order-events) ); } } Component public class RemoteEventReceiver { private final ApplicationEventPublisher publisher; private final ObjectMapper objectMapper; // 收到 Redis Stream 消息后触发示意生产代码需要 MessageListener public void onRemoteMessage(String payload) throws JsonProcessingException { OrderCreatedEvent event objectMapper.readValue(payload, OrderCreatedEvent.class); publisher.publishEvent(event); // 重新注入本地容器 } }看着挺合理。但线上环境和本地最大的不同是调用会跨进程、跨网络甚至跨机器、跨机房原本被 JVM 屏蔽掉的各种故障被重新引入。消息可能丢失消息可能重复消息可能乱序服务可能暂时不可用事件结构可能升级不兼容。这些距离带来的问题Spring Event 一个都没有解决。于是“看似好用”很快会被现实打脸。2.2 四个致命伤每一个都让人想删代码第一个致命伤是事件对象的序列化与版本兼容问题。本地事件传递的是对象引用不存在序列化问题一旦走向远程就必须把事件变成字节流。这时候你会被迫面对几个难题业务事件类上有没有实现Serializable时间字段用的是LocalDateTime还是Date事件类里有没有匿名内部类、Lambda、代理对象两个服务是否依赖完全一致的事件类定义只要事件结构里多了一个字段、少了一个字段旧消费者可能直接反序列化失败。很多团队为了省事把事件类放在一个 common 包里结果每次升级事件类都要同时发布两个服务等于制造了 API 层面的强耦合。第二个致命伤是丢消息与重复消费几乎不可避免。本地的publishEvent只要进程不崩事件一定被同进程的监听器收到远程传输则要面对网络超时、消费端宕机、消费者组再均衡、消息中间件自身故障。Redis Pub/Sub 模式下订阅端掉线期间的消息直接丢失Kafka 默认 at-least-once 语义可能重试多次RabbitMQ 手动确认机制也容易因为处理超时导致重复投递。Spring Event 本身不提供重试、确认、死信、幂等能力这些全都要你自己实现。第三个致命伤是分布式事务边界完全失效。本地用TransactionalEventListener(phase AFTER_COMMIT)可以把事件发送时机绑定到事务提交之后但事件一旦经过远程服务再被发布它处于另一个 JVM、另一个事务上下文中。如果业务要求“订单创建成功后通知积分服务两边要么都成功要么都不成功”远程 Spring Event 根本做不到。上游事务还没提交远端可能已经消费了上游回滚了远端已经发消息了。要保住最终一致性就得另做方案。第四个致命伤是顺序性和可观测性丢失。本地同步事件天然保证顺序发布端和监听端在同一个线程栈上日志天然关联远程场景下多实例并发消费意味着消息先后顺序完全不可控。再加上链路 IDTraceId如果不手动放入事件体排查问题时根本无法把一个远程事件跟源头请求串起来。生产环境一旦出问题你面对的是几十条日志和几个“不知道谁先谁后”的消息定位效率非常低。这四个问题组合到一起远不是“多写一层桥接”能解决的。所以很多团队在尝试过一轮之后都会得出同一个结论这不是在复用 Spring Event这是在用最差的方式实现一个消息队列。3. 别再拿事件总线当消息队列3.1 两个本质不同的模型要理解“远程 Spring Event 为什么没人再用”核心还是要分清两个模型。Spring Event 是一条进程内的内存总线本质是方法回调发布者调用监听器监听器在同一个 JVM 中执行没有独立存储、没有消费者组概念、没有重新投递机制。而消息队列MQ是跨进程的消息管道有持久化存储、有消费位点、有确认和重试机制天然面向分布式环境。打个比方Spring Event 像公司内部电话拨个分机就能通话速度快、线路稳定、挂了能直接找到人跨服务的消息队列像快递物流有面单、有仓库、有签收流程可以保证快件在路上丢了会补发送错了能有据可查。你不能指望公司内部电话把包裹送到另一个城市就像你不能指望 Spring Event 把事件可靠送到另一个服务。有人会问我不就是用 Redis 转发了一下吗严格来说转发了之后它就不再是“Spring Event”了而是“用 Redis 实现的事件驱动消息传递”。名字还叫 Event但你必须按消息系统的思路去治理它。否则你就会陷入一个尴尬境地Spring Event 的 API 简化了你的表达但底层没有任何分布式可靠性支撑出了问题只能干瞪眼。3.2 更成熟的替代方案本地消息表、Outbox、事务消息、消息中间件既然目标是从本地走向远程业界早就有一批成熟的“事件驱动 分布式可靠”方案不必执着于复刻 Spring Event 的 API。最经典的是本地消息表。业务数据写入数据库时在同一事务里写一条事件记录随后由一个定时任务扫描并发送事件到 MQ收到成功确认后再把本地记录标记为已发送。这种方案的优点是没有分布式事务也能保证“本地数据更新”和“事件落库”强一致缺点是定时器和消息表会增加一点代码且存在一定延迟。更现代的是 Outbox 模式本质是本地消息表的变体。把事件写入数据库中的 outbox 表通过 CDC 工具如 Debezium监听数据库 binlog将事件变更实时消费并投递到 Kafka。这样业务服务完全不用关心怎么发消息缺点是引入 CDC 组件运维复杂度会上升。如果只是单条消息要保证事务可以使用 RocketMQ 事务消息或 Kafka 事务配合幂等消费。事务消息先发送 half 消息业务事务结束后再提交或回滚中间件保证只有 commit 后的消息才被消费者看到。这套方案很成熟但对中间件版本和配置有一定要求。至于消息中间件本身Kafka、RocketMQ、RabbitMQ 都可以承担跨服务事件分发。选择标准很简单需要高吞吐和事件回放选 Kafka需要事务消息和低延迟选 RocketMQ希望轻量易维护选 RabbitMQ。还有一层更接近 Spring 生态的抽象Spring Cloud Stream它可以把 Kafka、RabbitMQ 等绑定成统一的StreamListener从开发体验上尽量贴近事件监听但底层依然是可靠的消息管道而不是 JVM 内回调。为什么要绕这么大一圈而不直接用“远程 Spring Event”因为可靠性是分布式系统的刚需。你当然可以自己写重试、幂等、死信、顺序控制但这些恰恰是消息中间件花了很多年才做好的能力。重复造轮子的成本远高于引入一个标准 MQ。4. 如果必须从“本地”走向“远程”怎么设计才不踩坑4.1 推荐路线业务层用 Spring Event传输层交给消息管道在一些落地场景里尤其是老系统改造业务代码已经大量用了 Spring Event想一步到位全部替换成 MQ 改动太大。我的建议是分层设计而不是把 Spring Event 和远程传输混在一起。核心思想是业务代码只依赖ApplicationEventPublisher跨服务的传输逻辑全部收敛在一个独立的事件桥组件里。具体的实现步骤可以拆成五步第一步定义统一的远程事件信封RemoteEventEnvelope至少包含这些字段事件唯一 IDeventId、事件类型type、事件源服务source、事件体payload、发生时间occurredAt、链路 IDtraceId。其中eventId和traceId是两个绝对不能省略的字段前者用于幂等后者用于排查。第二步业务代码照常发布本地事件例如publisher.publishEvent(new OrderCreatedEvent(orderId, userId))。这个事件只负责表达“订单已创建”的业务事实不关心谁消费、怎么消费。第三步在服务内部定义一个本地监听器专门监听所有需要跨服务的事件。它拿到OrderCreatedEvent后把业务事件转成RemoteEventEnvelope序列化成 JSON然后通过消息管道发送。Component public class RemoteEventBridge { private final KafkaTemplateString, String kafkaTemplate; private final ObjectMapper objectMapper; EventListener public void forwardOrderCreatedEvent(OrderCreatedEvent event) throws JsonProcessingException { RemoteEventEnvelope envelope RemoteEventEnvelope.builder() .eventId(IdUtil.createEventId()) .type(OrderCreatedEvent) .source(order-service) .payload(objectMapper.writeValueAsString(event)) .occurredAt(Instant.now()) .traceId(TraceUtil.currentTraceId()) .build(); kafkaTemplate.send(order-events, event.getOrderId(), objectMapper.writeValueAsString(envelope)); } }第四步远程服务消费到信封后解析type把payload反序列化成对应的领域事件再调用applicationEventPublisher.publishEvent(domainEvent)。到达这一步远程服务内部依然可以使用熟悉的EventListener来编写业务逻辑。第五步也是最容易被忽略的必须在消费者端做幂等。因为消息管道普遍存在“至少一次”的可能同一封事件可能在重启后重复投递。用eventId查一下去重表已经处理过就直接返回否则才执行本地事件发布。幂等判断要放在业务副作用发生之前不能等监听器执行完再判断。4.2 参数选择与避坑清单这套分层设计能跑通但有十几个细节踩过坑才能写全。我整理出最关键的几条。事件体序列化不要用 JDK 自带的ObjectOutputStream。Java 原生序列化性能差、格式乱而且强依赖类定义和serialVersionUID稍微改动就报InvalidClassException。选择 JSONJackson最稳妥配合事件类型type字段可以做到一定程度的向后兼容。如果字段删除了老消费者可以直接忽略。消息体大小要克制。不要把整个订单对象、用户快照、文件内容塞进事件体。事件应该是一个“已发生事实的轻量描述”对方需要的数据可以只传 ID由下游按需查询。否则网络包和消费端内存都会被拖垮。生产端发送要监控发送结果。Kafka 的send方法是异步的如果只是调用一下就完事发送失败的异常是感知不到的。要对发送结果回调做处理必要时记录一条日志或者落入待重发表。消费端要记得开启手动提交不要让自动提交把位点推到前面导致消息丢失。TransactionalEventListener的使用位置要明确。本地业务事件里可以放心用它保证事务提交后转发远程消费端再发布本地事件时仍然可以使用TransactionalEventListener但这时的“事务”是消费方法的事务而不是上游订单事务。如果希望“上游提交后再通知下游”应该在订单服务的桥接监听器上使用AFTER_COMMIT而不是把希望寄托在远端。链路 ID 要放进信封。现在很多团队用 SkyWalking、Micrometer Tracing 做链路追踪但跨服务的事件走消息管道时TraceId 默认不会自动跨消息传递。发布端生成事件时用TraceId填充信封消费端消费后把traceId重新放到当前线程上下文中整个处理链路才能串起来。重试要有上限和退避策略不能无脑单线程死循环。最终失败的消息要进死信队列或者死信 Topic保留原始信封方便人工补偿。5. 实际踩坑记录与排查思路5.1 四个典型的线上问题我们一个一个踩过第一个坑是 Redis Pub/Sub 丢消息。最开始我在两个服务之间用RedisTemplate.convertAndSend转发事件线下怎么测都通过。上线后早上高峰期出现用户付了款但没到账的情况一查存储服务在消息推送间隙刚好发生了重启而且还遇到过 Redis 主从切换期间订阅关系重建消息就直接丢了。Pub/Sub 本质是“发后即焚”没有持久化这种丢法完全正常。后来把 Redis Stream 的 ConsumerGroup 和 pending list 用上才把消失的消息捞回来。第二个坑是反序列化类型对不上。订单服务在今天的新版本里给OrderCreatedEvent增加了一个channel字段积分服务由于发布流程滞后用的还是旧版 common 包。消息到达积分服务后Jackson 反序列化旧版类时直接抛UnrecognizedPropertyException积分一直没有发放。这个问题促使我们把事件类型版本号带在信封里并约定新增字段只允许可空类型避免旧消费者被新字段砸死。第三个坑是重复消费导致重复发奖。某个活动服务消费事件后要给用户发积分因为消费者端处理超时Kafka 自动把消费位点往前拉了同一事件被消费了两次。由于没有幂等发奖逻辑执行了两遍。这是一个非常朴实也最容易忽略的问题。加入eventId幂等判断后重复事件才被拦截住。第四个坑是发事件时机不对。之前在一个交易服务里我先调用publishEvent再更新数据库结果下游监听器读到的是旧库存造成超卖。定位后发现是顺序和事务问题改成数据库更新完成、事务提交后再通过TransactionalEventListener(AFTER_COMMIT)转发才最终正常。5.2 选择建议速查表本地 Spring Event 还是远程 MQ场景推荐方案原因单体应用内部业务解耦Spring Event进程内调用无序列化和网络开销单体应用内异步处理耗时任务Spring Event Async简单、可控跨服务广播、削峰填谷Kafka / RocketMQ持久化、重试、消费者组跨服务且需要事务一致Outbox / 事务消息保证本地事务和消息发送一致需要按业务维度保证顺序Kafka 分区键同 key 路由到同一分区缓存失效、清理类轻量通知Redis Stream / 普通 MQ无需复杂可靠性时降低运维成本已有大量 Spring Event 代码、想渐进改造分层事件桥 Kafka / Stream保留业务表达传输层替换这张表在实际项目里基本能覆盖绝大多数场景。你会发现列出的方案里没有“远程 Spring Event”因为它不属于任何一格。你可以把它理解成一个过渡态用来解决“老代码迁移成本高”的问题但它不是长期方案。后期如果团队切换消息中间件只需要替换桥接层里那一个KafkaTemplate或RedisTemplate业务监听器完全不动。这是一条真正能从“本地 Spring Event”平滑过渡到“远程消息驱动”的路径。我的经验是在本地Spring Event 是一个非常趁手的内务工具在远程它是被过度寄予厚望的替身。真正要跨服务发送事件不如直接承认“我需要一个消息队列”然后选一个成熟的中间件别让业务代码为一个不落地的事件总线买单。如果你现在正准备把一个本地事件改成远程通知请先回答三个问题对方服务暂停时消息能不能等同一事件重复消费能不能扛住跨服务的先后顺序丢失了要不要紧这三个问题只要有一个答不上来就应该放弃“远程 Spring Event”这种拼装方案踏实选一个 MQ。
返回列表