十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Flower三大消息处理模式详解:消息分叉、消息聚合与消息回复如何构建复杂业务流程

Flower三大消息处理模式详解:消息分叉、消息聚合与消息回复如何构建复杂业务流程 Flower三大消息处理模式详解消息分叉、消息聚合与消息回复如何构建复杂业务流程【免费下载链接】flower反应式微服务框架Flower项目地址: https://gitcode.com/gh_mirrors/flow/flowerFlower 是一个构建在 Akka 之上的反应式微服务框架其核心思想是消息驱动每个 Service 完成一个细粒度业务功能前一个 Service 的返回值会被框架封装成消息自动投递给后续 Service。对于新手来说理解 Flower 的消息处理模式消息分叉、消息聚合、消息回复是掌握这套反应式微服务框架的关键。本文用一套真实示例带你快速吃透这三大模式如何组合出复杂的业务流程。先认识 Flower 的服务编排流程即消息管道在 Flower 中流程是核心设计目标。开发者只需把多个 Service 按业务流程连线框架就会加载成消息的处理流通道。整体架构上网关层的 Controller 绑定流程注册中心负责服务编排各 Flower 容器中的 Service 通过消息接力完成业务并返回结果框架内部ServiceFlow定义了谁的后继是谁ServiceActor负责把消息按后继列表逐个投递。这套消息投递的完整时序如下看懂这张图三大消息处理模式其实只是后继列表的不同玩法。消息分叉一条消息喂给多个并行服务消息分叉指一个服务输出的消息分发给 1 个或多个其他服务分两种形式全部分发消息发给全部后继服务并行执行。典型场景是用户注册成功后写库、发激活邮件、发通知短信、同步关联产品4 个任务同时开跑。条件分发根据消息内容只发给其中一个后继服务比如贷款申请按信用等级选择不同审批服务。条件分发有三种实现方式按消息泛型类型匹配后继服务声明ServiceMessageB就只接收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)注解标记自己为聚合节点用ListObject接收三路分叉的聚合结果FlowerService(type FlowerType.AGGREGATE) public class ServiceReceiveA implements ServiceListObject, Integer { Override public Integer process(ListObject 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配置为内置AggregateServiceservice4从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.syncCallServiceFlower 的精髓在于分叉、聚合、回复都只是对消息 后继列表的不同配置服务代码保持无阻塞、可复用。掌握这三大消息处理模式你就能用 Flower 编排出任意复杂度的反应式业务流程。【免费下载链接】flower反应式微服务框架Flower项目地址: https://gitcode.com/gh_mirrors/flow/flower创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表