Flower三大消息处理模式详解:消息分叉、消息聚合与消息回复如何构建复杂业务流程
【免费下载链接】flower反应式微服务框架Flower项目地址: https://gitcode.com/gh_mirrors/flow/flower
Flower 是一个构建在 Akka 之上的反应式微服务框架,其核心思想是消息驱动:每个 Service 完成一个细粒度业务功能,前一个 Service 的返回值会被框架封装成消息,自动投递给后续 Service。对于新手来说,理解 Flower 的消息处理模式(消息分叉、消息聚合、消息回复)是掌握这套反应式微服务框架的关键。本文用一套真实示例,带你快速吃透这三大模式如何组合出复杂的业务流程。
先认识 Flower 的服务编排:流程即消息管道
在 Flower 中,流程是核心设计目标。开发者只需把多个 Service 按业务流程"连线",框架就会加载成消息的处理流通道。整体架构上,网关层的 Controller 绑定流程,注册中心负责服务编排,各 Flower 容器中的 Service 通过消息接力完成业务并返回结果:
框架内部,ServiceFlow定义了"谁的后继是谁",ServiceActor负责把消息按后继列表逐个投递。这套消息投递的完整时序如下:
看懂这张图,三大消息处理模式其实只是"后继列表"的不同玩法。
消息分叉:一条消息喂给多个并行服务
消息分叉指一个服务输出的消息分发给 1 个或多个其他服务,分两种形式:
- 全部分发:消息发给全部后继服务,并行执行。典型场景是用户注册成功后,"写库、发激活邮件、发通知短信、同步关联产品"4 个任务同时开跑。
- 条件分发:根据消息内容只发给其中一个后继服务,比如贷款申请按信用等级选择不同审批服务。
条件分发有三种实现方式:按消息泛型类型匹配(后继服务声明Service<MessageB>就只接收MessageB)、消息实现Condition接口指定后继服务 id、以及使用框架内置的ConditionService做通用分发。
一个完整示例见 flower.sample 的 aggregate 包,AggregateController用buildFlow声明了"两次分叉":
// 第一个分叉:ServiceBegin 的消息发给 A1、A2、A3 三个服务并行处理 getServiceFlow().buildFlow(ServiceBegin.class, ServiceForkA1.class); getServiceFlow().buildFlow(ServiceBegin.class, ServiceForkA2.class); getServiceFlow().buildFlow(ServiceBegin.class, ServiceForkA3.class); // 三个分叉结果再汇聚到 ServiceReceiveA getServiceFlow().buildFlow(ServiceForkA1.class, ServiceReceiveA.class); getServiceFlow().buildFlow(ServiceForkA2.class, ServiceReceiveA.class); getServiceFlow().buildFlow(ServiceForkA3.class, ServiceReceiveA.class);可以看到,分叉的本质只是在流程中给同一个前序服务配置了多个后继服务,服务代码本身零改动。
消息聚合:把并行结果收拢成一个集合
分叉的目的通常是并行提速,而消息聚合则负责把多个并行服务产生的结果收拢起来,交给后续服务统一处理。
框架内置了聚合服务AggregateService,它会把多路消息封装成一个Set(或List)返回。示例中ServiceReceiveA使用@FlowerService(type = FlowerType.AGGREGATE)注解标记自己为聚合节点,用List<Object>接收三路分叉的聚合结果:
@FlowerService(type = FlowerType.AGGREGATE) public class ServiceReceiveA implements Service<List<Object>, Integer> { @Override public Integer process(List<Object> message, ServiceContext context) { // message 是 A1、A2、A3 三个并行服务结果的集合 Integer sum = 0; for (Object obj : message) { if (obj instanceof Integer) { sum += (Integer) obj; } } return sum; // 求和后继续触发 B 分叉 } }文件见flower.sample/src/main/java/com/ly/train/flower/sample/aggregate/service/ServiceReceiveA.java。聚合之后流程并没有结束,ServiceReceiveA的返回值又触发了第二轮分叉(B1、B2),最终由同样标注FlowerType.AGGREGATE的ServiceReceiveAB再次聚合收口——分叉和聚合可以无限嵌套组合,这就是复杂业务流程的构建方式。
更简单的聚合用法见 textflow 示例:service2 -> service5、service3 -> service5两条支线汇聚,service5配置为内置AggregateService,service4从Set中取出各结果分别处理。
消息回复:不阻塞也能拿到最终结果
Flower 中消息全部异步处理,服务之间不互相阻塞等待,这正是低耦合、无阻塞、高并发的来源。流程调用者发出请求后无需等待,可以继续处理下一个请求。
但很多场景调用方确实需要拿到最终处理结果,这时就用到消息回复模式:通过ServiceFacade的syncCallService发起调用,框架会把流程的最终结果消息回复给调用者:
Message2 m2 = new Message2(10, "Zhihui"); Message1 m1 = new Message1(); m1.setM2(m2); // 消息回复:阻塞式获得整条流程的最终处理结果 System.out.println("返回结果" + serviceFacade.syncCallService("sample", m1));注意区分两种调用姿势:asyncCallService发出即返回,适合高吞吐场景;syncCallService等待最终消息回复,适合需要同步取值的接口入口。Web 场景中,流程内的服务也可以直接通过HttpService把结果推送给用户端,发起请求的 Servlet 完全不必等待——两种模式按业务需要自由搭配。
三大模式组合:5 分钟看懂一个完整流程
把示例串起来看(flower.sample/src/main/java/com/ly/train/flower/sample/aggregate/AggregateApplication.java),完整流程是:
┌─ ForkA1 ─┐ Begin ─分叉─┼─ ForkA2 ─┼─→ ReceiveA(聚合)─分叉─┬─ ForkB1 ─┬→ ReceiveAB(聚合收口) └─ ForkA3 ─┘ └─ ForkB2 ─┘- 分叉:
ServiceBegin的一条消息并行触发 A1/A2/A3 三个任务; - 聚合:
ServiceReceiveA用FlowerType.AGGREGATE收拢三路结果并求和; - 再分叉、再聚合:求和结果触发 B 分叉,
ServiceReceiveAB二次聚合; - 回复:Web 请求通过
FlowerController绑定这条流程(@Flower(value = "aggregate", flowNumber = 6)声明通道数),最终结果回复给调用方。
想动手体验,可克隆仓库https://gitcode.com/gh_mirrors/flow/flower,运行 aggregate 示例并访问/test/aggregate/{id}接口,观察控制台输出的分叉与聚合日志。
总结
| 模式 | 解决的问题 | 关键实现 |
|---|---|---|
| 消息分叉 | 并行执行多个子任务 | buildFlow配置多个后继 /Condition条件分发 |
| 消息聚合 | 收拢并行结果 | 内置AggregateService/FlowerType.AGGREGATE |
| 消息回复 | 异步体系中同步取值 | ServiceFacade.syncCallService |
Flower 的精髓在于:分叉、聚合、回复都只是对"消息 + 后继列表"的不同配置,服务代码保持无阻塞、可复用。掌握这三大消息处理模式,你就能用 Flower 编排出任意复杂度的反应式业务流程。
【免费下载链接】flower反应式微服务框架Flower项目地址: https://gitcode.com/gh_mirrors/flow/flower
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考