ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spring Messaging抽象层:微服务消息驱动架构的解耦之道与踩坑指南

Spring Messaging抽象层:微服务消息驱动架构的解耦之道与踩坑指南 消息驱动架构现在是微服务绕不开的话题但很多人一上来就扎进 Kafka、RabbitMQ 的客户端 API 里写出来的代码和中间件强耦合换个 MQ 几乎等于重写。Spring Messaging 这套抽象层其实才是真正值得先搞明白的东西它把消息本身、消息通道、消息处理器这些概念统一了起来不管你底层用哪种消息中间件上层代码的骨架都可以保持稳定。这篇内容我会结合自己实际做过的项目把 Spring Messaging 的模型拆开讲清楚再延伸到和 Spring Integration、Spring Cloud Stream 的关系最后聊几个高频踩坑点比如重复消费、顺序性、事务消息到底该在哪一层解决。1. 消息驱动架构里Spring Messaging究竟在解决什么问题1.1 不是又一个消息中间件而是一层抽象规范先纠正一个常见的误解Spring Messaging 并不是拿来替代 Kafka、RabbitMQ 或 RocketMQ 的东西它本身不负责消息的存储、投递、持久化也不管你消息是走 TCP 还是走 HTTP。它的定位是 Java 领域里一套面向消息的统一编程模型解决的是业务代码和具体消息中间件耦合太深这一类问题。我见过不少项目Service 层里直接注入 RocketMQ 的 producer然后调用 send 方法消息体还强依赖某个 MQ 的 Message 类型。这样做短期没问题但一旦要换中间件或者同一个项目里需要对接两种 MQ业务代码就会变得很难看。Spring Messaging 的出发点就是把这些细节挡在业务外面让你面向 Message、MessageChannel、MessageHandler 这几个抽象编程具体的传输协议和中间件适配由底层实现去处理。这个思路其实和 Spring 家族一贯的作风一致先定义接口和约定再通过不同实现去适配各种场景。Spring Messaging 对应的是 spring-messaging 模块Spring WebSocket、Spring Integration、Spring Cloud Stream 都以它作为基础。倒过来说也一样如果你理解了 Spring Messaging 的模型后面接触那些上层框架时会顺畅很多。1.2 核心概念的映射关系为了后面讲起来不绕先把 Spring Messaging 里几个顶级抽象和日常开发里的概念做一个对照Spring Messaging 抽象你可以理解成核心职责Message一封信携带消息体和消息头MessageHeaders信封上的收件人、时间戳等信息描述消息的元数据MessageChannel邮局里的传送管道负责消息的传递MessageHandler处理信件的人负责处理消息内容MessagingTemplate帮你寄信的服务台提供发送消息的便捷API这个对照不绝对精确但足够帮你建立第一印象。实际写代码的时候我们通常不直接 new 一个 MessageChannel 然后手动 send更多是借助注解或者模板类去完成消息的收发。不过理解这些底层概念对排查问题和阅读框架源码都很有帮助。1.3 什么时候你会真正用到它Spring Messaging 最直接的落地场景是 WebSocket 消息推送。你在后端往某个用户或者某个会话推送一条消息底层走的是 WebSocket 协议但在 Spring 的世界里你打交道的对象就是 Message 和 MessageChannel。另外如果项目里引入了 Spring Integration那整套企业集成模式EIP的信道、路由器、过滤器全都是基于 Spring Messaging 构建的。再往上走Spring Cloud Stream 的输入输出通道本质上也是一条条 MessageChannel。所以可以先记住一个结论Spring Messaging 不是让你直接拿来发消息给 Kafka 的它是所有 Spring 消息体系的底层契约。2. Message、MessageHeaders、Channel那套被说烂了却很少讲透的核心API2.1 Message 为什么长这样直接看接口定义会比较容易理解。Spring Messaging 的 Message 长这样public interface MessageT { T getPayload(); MessageHeaders getHeaders(); }就这么简单。一个 payload一组 headers。payload 是真正的业务数据headers 是元数据。你可能会觉得这也太简陋了但仔细想一下任何消息系统最终都逃不开这两个东西你发给别人的内容以及描述这段内容的附加信息。Spring 的 MessageHeaders 继承自 MapString, Object但它的实现比普通 map 严谨一点内部维护了一个固定的 header 键集合比如 id 和 timestamp 在创建时就会被自动加上而且不允许修改。这保证了每条消息在整个传递过程中都具备全局唯一的 id 和时间戳方便做链路追踪也方便消息消费方做去重。实际创建消息也不用手动去 new 一个实现类直接用 GenericMessage 或者 MessageBuilderMessageString message MessageBuilder.withPayload(hello) .setHeader(appId, order-service) .build();或者用 GenericMessage 也行MessageString message new GenericMessage(hello); GenericMessage 内部会帮你生成 id 和 timestamp 这两个系统 header所以你可以直接把它扔到 MessageChannel 里去。这就是 Spring Messaging 的第一个实用体验消息本身的创建成本很低不需要依赖任何中间件客户端。 ### 2.2 MessageChannel只管往管道里塞不负责你想要的结果 MessageChannel 是消息传递的管道抽象接口定义也极其克制 java public interface MessageChannel { boolean send(Message? message); boolean send(Message? message, long timeout); }send 方法返回 boolean表示消息是否被成功接收注意是接收而不是处理。Channel 不会保证消息被业务逻辑成功消费它只保证我把消息交给了下一个节点。这是一个很重要的心智模型很多人在排查消息丢失问题时会把责任归到 Channel 上其实 Channel 只负责传递业务上的异常需要靠 MessageHandler 去捕获。MessageChannel 有两个常见子接口SubscribableChannel可以注册多个 MessageHandler 订阅者有点类似发布-订阅模型。PollableChannel支持主动拉取消息调用 receive() 方法从管道中取消息。这两个子接口的区别对应了两种不同的消费模式事件驱动和轮询驱动。前者适合实时性要求高的场景后者适合需要批量拉取、控制消费节奏的场景。2.3 MessageHandler消息到了之后由你决定怎么处理MessageHandler 接口更简单public interface MessageHandler { void handleMessage(Message? message) throws MessagingException; }很多刚接触的同学会问这不就是一个消费者吗和 MQ 的 consumer 有什么区别区别在于 MessageHandler 不绑定任何具体的消息中间件。它就是一个纯粹的业务处理单元输入是 Message输出由你决定。你可以把它注册到 SubscribableChannel 上让它持续接收消息也可以把它封装成某个消息中间件的 listener接住 Kafka 或 RabbitMQ 的消息再转交给它。我自己在做项目时很喜欢把核心业务逻辑全部放在 MessageHandler 里然后外部再通过一层薄薄的 adapter 对接具体中间件。这样中间件相关的代码只出现在项目边缘核心代码始终保持纯粹。2.4 一套最简可运行的示例把上面的概念串一下一个最简的基于 Spring Messaging 的消息收发程序可以写成这样Configuration public class MessagingConfig { Bean public SubscribableChannel myChannel() { return new DirectChannel(); } Bean public MessageHandler myHandler() { return message - System.out.println(收到消息: message.getPayload()); } }然后在使用方注入 channel把消息发进去Component public class MessageSender { private final SubscribableChannel myChannel; public MessageSender(SubscribableChannel myChannel) { this.myChannel myChannel; this.myChannel.subscribe(message - System.out.println(订阅者A: message.getPayload())); } public void send(String content) { myChannel.send(new GenericMessage(content)); } }DirectChannel 是 SubscribableChannel 的默认实现它是同步执行的消息发送出去后会直接在当前线程调用订阅者。这个特性后面讲线程模型时还会提。提示千万别在 DirectChannel 的订阅者里做耗时的 IO 操作它会阻塞发送方线程。需要异步处理时应该换成 ExecutorSubscribableChannel或者交给上层消息中间件去异步投递。3. 注解消息映射MessageMapping和SendTo如何做到像写Controller一样处理消息3.1 从 Spring MVC 到 Spring Messaging 的注解思想迁移如果你写过 Spring MVC那对 RequestMapping、ResponseBody 这套注解应该很熟。Spring Messaging 在 WebSocket 场景下提供了一套几乎对称的注解模型让你可以把 WebSocket 消息直接映射到 Controller 的方法上。这一套注解模型和 spring-messaging 模块有关它定义的 MessageMapping、SendTo、Payload、Header 等注解使得处理一条 WebSocket 消息的代码看起来和处理一条 HTTP 请求的代码没有太大差别。对于团队来说这种一致性降低了学习成本也让代码结构更统一。在 WebSocket 的 STOMP 协议场景下前端会发送类似这样的帧{destination: /app/chat, payload: 你好}后端可以用 MessageMapping 直接接收Controller public class ChatController { MessageMapping(/chat) SendTo(/topic/greetings) public Greeting chat(Payload String message) { return new Greeting(收到: message); } }当一个用户通过 WebSocket 往/app/chat发送消息时Spring 会自动把它转换成 Message解析出 payload调到这个方法上再把返回值封装成新的 Message 发送到/topic/greetings这个广播地址。3.2 参数绑定不只是取 payload 这么简单MessageMapping 的方法可以接收的参数类型比很多人想象中要丰富。除了 Payload 可以绑定消息体之外Header 可以绑定消息头DestinationVariable 可以绑定路径模板变量。举例来说MessageMapping(/topic/{roomId}) SendTo(/topic/room/{roomId}) public ChatMessage sendToRoom(DestinationVariable String roomId, Payload ChatMessage payload) { return payload; }这个方法里roomId 从目标地址里解析出来payload 从消息体里解析出来。这样的代码写起来非常清爽不再需要手动去解析消息帧的格式。如果消息体是 JSONSpring 会利用配置好的 MessageConverter 自动做反序列化把 JSON 映射成 ChatMessage 对象。常见的就是 Jackson 的 converterspring-messaging 已经帮你集成好了。如果消息体本身就是 String那直接声明 String 类型的 payload 参数就行。3.3 如何在 Spring Boot 里把 WebSocket 消息通道配出来注解只是入口要让整条链路真正跑起来还得配置 WebSocket 的消息代理。这里直接给一个我在项目里常用的配置Configuration EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { Override public void configureMessageBroker(MessageBrokerRegistry registry) { registry.enableSimpleBroker(/topic, /queue); registry.setApplicationDestinationPrefixes(/app); } Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint(/ws).setAllowedOriginPatterns(*).withSockJS(); } }这组配置的含义是客户端可以通过 SockJS 连接到/ws端点服务端会在内存里启一个简单的消息代理负责把消息转发给订阅了/topic或/queue的客户端而客户端发送消息时目标地址以/app开头会交给 MessageMapping 注解的方法处理。这里我特别想强调 enableSimpleBroker 的坑。SimpleBroker 是一个内存内的简易消息代理只能跑在单机环境下。如果部署了多个实例客户端 A 连接到实例 1客户端 B 连接到实例 2那么实例 1 上广播的消息实例 2 根本收不到。生产环境必须引入外部消息中间件做 STOMP broker relay比如 RabbitMQ。这也是很多人从 Demo 走向生产时第一次遇到消息丢了的常见原因。3.4 运维视角为什么注解消息映射适合做实时推送类接口我负责过一个实时监控大屏项目前端需要通过 WebSocket 接收后端的指标变化。后端其实就是用 MessageMapping 接收前端订阅指令再通过 SimpMessagingTemplate 主动推送数据。SimpMessagingTemplate 是 Spring Messaging 提供的又一个实用工具可以在任意地方向指定用户或指定地址推送消息simpMessagingTemplate.convertAndSend(/topic/metrics, metricData);也可以向单个用户推送simpMessagingTemplate.convertAndSendToUser(username, /queue/message, data);convertAndSendToUser 和 /queue 前缀配合可以实现点对点推送。这种接口模式比 HTTP 长轮询的体验好很多而且借助 Spring Messaging 的抽象业务代码里不需要出现任何 WebSocket 原生 API。4. 路由、过滤、转换消息管道里的三个关键角色4.1 消息不能只靠 Channel 直通管道里还需要处理器链真实项目里的消息流往往不是一条直线走到底。相同类型的消息可能需要根据消息头里的字段被分发到不同的处理节点有些消息可能不符合规则需要被丢弃有些消息在进入业务逻辑前需要做格式转换。Spring Messaging 虽然模块本身没有像 Spring Integration 那样完整的 EIP 支持但它定义的 MessageChannel 和 MessageHandler 模式足以支撑这些规则。有一种很实用的设计模式一个 MessageChannel 作为入口挂载多个 MessageHandler每个 MessageHandler 内部判断自己是否应该处理这条消息如果符合条件就处理如果不符合就原样转发到下一条 Channel。这种链式结构让每个处理节点都可以保持单一职责。4.2 路由器从一条消息里决定下一步去哪这里我拿 Spring Integration 来做示例因为它把 Spring Messaging 的 Channel 模型用得很彻底。假设你需要根据消息里的 orderType 字段把订单消息路由到不同的处理器Bean public IntegrationFlow orderRouterFlow() { return IntegrationFlow.from(orderInputChannel) .OrderMessage, Stringroute(OrderMessage::getOrderType, mapping - mapping .subFlowMapping(NORMAL, sub - sub.handle(normalOrderHandler)) .subFlowMapping(GIFT, sub - sub.handle(giftOrderHandler))) .get(); }route 方法接收一个 SpEL 表达式或者一个 Function返回值用来做路由键。这个路由键决定消息进入哪个子流。这种写法的基础还是 MessageChannel 和 MessageHandler只是 Spring Integration 的 DSL 把细节封装得很舒服。4.3 消息过滤器不该处理的直接丢弃过滤器的逻辑更简单判断条件不满足就把消息过滤掉。在 Spring Integration 里可以这样写Bean public IntegrationFlow filteredFlow() { return IntegrationFlow.from(rawMessageChannel) .filter(Message::getHeaders, h - VIP.equals(h.get(userLevel))) .handle(vipPromotionHandler) .get(); }这段代码的含义很直白只有 userLevel 等于 VIP 的消息才会继续往下走其余消息在这一个节点就被拦截了。如果你用的是纯 Spring Messaging 而不引入 Spring Integration也可以自己写一个 CompositeMessageHandler 在内部做条件判断效果类似。4.4 消息转换器把外部消息翻译成内部对象消息转换在对接外部系统时特别常见。举个例子订单服务通过 Kafka 收到一条 JSON 格式的创建订单事件但你希望业务层只看到 OrderCreatedEvent 对象而不是一个原始的 JSON 字符串。消息转换器就是干这个的Bean public IntegrationFlow orderEventFlow() { return IntegrationFlow.from(orderEventChannel) .transform(JsonToObjectTransformer.class, spec - spec .transformer(new JsonToObjectTransformer(OrderCreatedEvent.class))) .handle(orderEventHandler) .get(); }转换器把 Message 的 payload 从一种形式变成另一种形式然后封装成新的 Message 继续往下传。这套机制在 Spring Messaging 里的底层就是 PayloadTypeConvertingChannel 或者 Transformer 处理器核心思想都是每个节点只做一件事做完以后把消息交给下一个节点。我个人的经验是路由、过滤、转换这三个操作尽量在消息进入业务逻辑之前完成。业务层的代码只需要面对已经清洗好的领域对象不需要关心消息是从哪个系统来的、格式是什么。5. Spring Messaging、Spring Integration、Spring Cloud Stream 的分层关系与选型边界5.1 三个框架不是竞争关系是层层递进很多同学对这三兄弟的关系一直有些模糊。这里我用一个比较通俗的方式解释Spring Messaging定义了消息、通道、处理器这些基础概念相当于一套通用的消息编程 API。Spring Integration在 Spring Messaging 之上实现了企业集成模式提供了丰富的消息通道、路由器、过滤器、转换器、网关等组件解决的是系统内部不同模块之间如何用消息协作的问题。Spring Cloud Stream在 Spring Messaging 之上进一步抽象了与消息中间件的绑定逻辑解决的是微服务之间如何通过消息队列通信的问题。Spring Cloud Stream 应用无论底层绑定的是 Kafka 还是 RabbitMQ编程模型都用 Input 和 Output 注解把一个 Channel 和一个 Java 方法绑定起来。现代版本推荐用函数式编程模型直接把 Supplier、Function、Consumer 作为消息端点。5.2 实际项目中怎么选我给一个选型建议这个建议来自我踩过的一些弯路场景推荐方案原因单个应用内部模块间异步解耦Spring Events Spring Messaging轻量不需要引入中间件需要复杂消息流编排、多节点转化Spring IntegrationEIP 模式完整DSL 表达力强微服务之间通过 MQ 通信Spring Cloud Stream绑定抽象好切换 MQ 成本低WebSocket 实时推送Spring Messaging STOMP天然支持生态成熟这里要注意避免一种情况项目只是简单发几条消息到 Kafka但为了规范化强行把 Spring Cloud Stream 全家桶引进来。框架的选择应该适配业务的复杂度而不是堆得越多越好。5.3 Spring Cloud Stream 在 Spring Messaging 之上还做了什么Spring Cloud Stream 的核心价值在于 Binder 机制。Binder 负责把外部消息中间件连接到一个 Channel 上。应用代码只需要声明一个输入通道和一个输出通道Bean public FunctionFluxString, FluxString process() { return input - input.map(String::toUpperCase) .doOnNext(System.out::println); }这个函数式接口表示一个处理链路从输入 Flux 接收消息做转换再输出到另一个 Flux。Spring Cloud Stream 会根据配置把输入 Flux 连接到某个 Kafka topic 或 RabbitMQ queue。你看这个过程中应用代码里没有出现任何 Kafka 或 RabbitMQ 的客户端类这就是在 Spring Messaging 之上做抽象的好处。如果直接裸用 Kafka 客户端也不是不行但你会发现自己写了不少重复的配置代码序列化器、消费者组、重试机制、offset 管理。Spring Cloud Stream 帮你把这一层封装收敛了换来的是写业务代码的时间。5.4 什么时候只靠 Spring Messaging 就够了反过来我也要泼一盆冷水。如果你的系统里根本没有使用外部消息中间件只是在单个应用内部做一些事件解耦那引入 Spring Integration 或者 Spring Cloud Stream 都属于过度设计。Spring Messaging 自带的 MessageChannel 和 MessageHandler 已经可以满足大部分需求配合 Spring Events 足以应付常见的场景。孤立地看 Spring Messaging你会发现它的 API 非常小但这恰恰是它的优势边界足够清晰不会干扰你现有的技术栈。6. 消息中间件场景下的常见陷阱重复消费、顺序性和事务消息该在哪里解决6.1 重复消费问题的根因与解法只要用了消息队列重复消费几乎必然会遇到。常见原因包括消费端处理完业务但还没来得及提交 offset进程挂了消息被重新投递或者消费端处理逻辑里做了重试但前一次其实已经成功了。Spring Messaging 的抽象层面并不负责去重它只是把消息从通道里取出来交给你的 MessageHandler。去重必须在业务代码里自己控制。最稳妥的做法是引入一张去重表或者利用 Redis 的 setnx 命令以消息的唯一 ID 做幂等标记。Spring Messaging 的 Message 自带 id 消息头这个 id 在消息创建时生成。如果你是通过 Spring Cloud Stream 消费 Kafka 消息可以在消息头里找到 kafka 相关的 offset 信息用topic partition offset拼接成一个业务幂等键。代码大致长这样public void handleMessage(MessageString message) { String dedupeKey message.getHeaders().get(kafka_receivedPartitionId) - message.getHeaders().get(kafka_offset); Boolean first stringRedisTemplate.opsForValue() .setIfAbsent(dedupeKey, 1, Duration.ofMinutes(10)); if (Boolean.TRUE.equals(first)) { process(message.getPayload()); } }setnx 操作如果返回 true说明这条消息第一次处理业务照常走如果返回 false说明已经被处理过直接丢弃。这里有一个细节值得注意幂等键一定要加过期时间否则去重表会无限膨胀。过期时间可以根据业务的容忍度来定比如允许 10 分钟内重复。6.2 消息顺序性Spring 不会帮你保证分区键才是关键消息顺序性的典型需求是同一个订单的创建、支付、完成事件必须按顺序被消费。如果并发消费就有可能出现支付事件先于创建事件到达业务层直接导致数据异常。Kafka 层面的解法是保证同一个 key 的消息进同一个分区同一个分区内的消息消费是严格有序的。生产者要指定消息 keyProducerRecordString, String record new ProducerRecord( order-events, orderId, eventJson);同一个 orderId 会落入同一分区。但是如果你的 Spring Cloud Stream 消费端并发度设置成大于分区数那多个线程会同时拉取同一个分区内的消息顺序仍可能被打乱。Spring Cloud Stream 里有一个配置项可以控制并发消费者数量要保证强顺序通常要把并发度降为 1spring.cloud.stream.bindings.input.consumer.concurrency1这行配置的含义是输入通道只有一个消费者线程在消费。这样做虽然损失了一部分吞吐量但能确保同一分区内的消息顺序不乱。如果你的业务对吞吐量和顺序性都有很高要求那就要在业务层设计更好的方案比如基于状态机的乱序消息聚合但这属于更高的复杂度了。Spring Messaging 自身不会去处理这些因为它只负责消息在应用内部流转不负责消息跨系统的投递语义。认清这一点排查问题就能少走很多弯路。6.3 事务消息抽象框架没有统一的 API别指望框架帮你搞定事务消息是 RocketMQ 提出的一种能力用于解决本地事务和消息发送不一致的问题。核心思路是先发一条半消息本地事务执行成功后 commit执行失败则 rollback消息中间件会定期回查事务状态。Spring Messaging 和 Spring Cloud Stream 都没有对事务消息做一套统一的标准 API。RocketMQ 的事务消息需要依赖 RocketMQ 自己的客户端类来实现TransactionMQProducer producer new TransactionMQProducer(group); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 执行本地事务 boolean success doLocalBusiness(); return success ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 回查事务结果 return LocalTransactionState.COMMIT_MESSAGE; } });如果你在 Spring Cloud Stream 里用 RocketMQ binder部分版本支持发送事务消息但配置方式和原生 API 不完全一致需要查看具体 binder 的文档。我的建议是如果事务消息是你系统的核心诉求不要试图用通用抽象盖住所有的差异直接把 RocketMQ 的客户端纳入基础设施层让业务代码通过一个封装好的接口来调用。6.4 消费端健壮性从 MessageHandler 到死信队列最后聊一个很多人会在生产环境遇到的状况消息处理失败了怎么办有些团队选择在 catch 里打日志然后继续消费这等于把消息丢了有些团队选择无限重试结果把下游系统拖垮。比较合理的做法是在 MessageHandler 里区分可重试异常和不可重试异常。可重试异常比如下游临时超时可以交给消息中间件的重试机制来处理不可重试异常比如消息体格式错误、业务数据不满足前置条件应该记录错误日志后把消息投入死信队列或者落库留待人工处理。Spring Cloud Stream 里可以为每个 binding 配置重试参数spring.cloud.stream.bindings.input.consumer.max-attempts3 spring.cloud.stream.bindings.input.consumer.back-off-initial-interval1000 spring.cloud.stream.bindings.input.consumer.back-off-max-interval3000三次重试仍然失败消息会进入 DLQDead Letter Queue之后由专门的补偿任务去处理。这样既不会丢消息也不会因为无脑重试把系统搞挂。我在实际处理这类问题时还有一个习惯MessageHandler 里显式捕获异常并打印消费的消息头和异常堆栈这样即便消息进了死信队列排查时也有足够上下文。日志里如果只有一条消费失败而没有消息 ID 和消息头信息排查效率会低很多。从 Spring Messaging 的 MessageHandler 到 Spring Integration 的过滤器再到 Spring Cloud Stream 的消费端配置整条链路看下来你会发现Spring 这套消息体系的设计核心始终是分层。消息如何传输、如何保证可靠性、如何重试这些都可以通过配置和组件去扩展而业务代码始终只面向统一的消息模型。我后来再看各类消息中间件的落地方案时都会先想想它属于哪一层该用哪一层的工具去解决想清楚之后代码写起来会清爽很多。
返回列表