news 2026/8/28 12:51:41

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flower三大消息处理模式详解:消息分叉、消息聚合与消息回复如何构建复杂业务流程

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 包,AggregateControllerbuildFlow声明了"两次分叉":

// 第一个分叉: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.AGGREGATEServiceReceiveAB再次聚合收口——分叉和聚合可以无限嵌套组合,这就是复杂业务流程的构建方式。

更简单的聚合用法见 textflow 示例:service2 -> service5service3 -> service5两条支线汇聚,service5配置为内置AggregateServiceservice4Set中取出各结果分别处理。

消息回复:不阻塞也能拿到最终结果

Flower 中消息全部异步处理,服务之间不互相阻塞等待,这正是低耦合、无阻塞、高并发的来源。流程调用者发出请求后无需等待,可以继续处理下一个请求。

但很多场景调用方确实需要拿到最终处理结果,这时就用到消息回复模式:通过ServiceFacadesyncCallService发起调用,框架会把流程的最终结果消息回复给调用者:

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 ─┘
  1. 分叉ServiceBegin的一条消息并行触发 A1/A2/A3 三个任务;
  2. 聚合ServiceReceiveAFlowerType.AGGREGATE收拢三路结果并求和;
  3. 再分叉、再聚合:求和结果触发 B 分叉,ServiceReceiveAB二次聚合;
  4. 回复: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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/28 12:50:17

架构图排障证据的保留

架构图排障证据的保留本文讨论“用精美架构图讲清复杂技术原理的方法&#xff1a;日志、指标、Trace 的可观测性落地”。这里的场景只用于说明方法&#xff0c;不对应某次真实线上事故&#xff1b;任何效果判断都应以项目自己的数据、配置和负载为准。 先定义问题和边界 “用精…

作者头像 李华
网站建设 2026/8/28 12:49:58

Browser-Use 文件下载自动化实战:2个参数跑通全链路

Browser-Use 文件下载自动化实战&#xff1a;2个参数跑通全链路 【免费下载链接】browser-use &#x1f310; Make websites accessible for AI agents. Automate tasks online with ease. 项目地址: https://gitcode.com/GitHub_Trending/br/browser-use 让 Agent 去 f…

作者头像 李华
网站建设 2026/8/28 12:47:21

数学建模竞赛实战指南:从团队分工到72小时极限攻关

1. 项目概述&#xff1a;一次从零到一的数模竞赛深度复盘 “2021全国数学建模大学”——这个标题&#xff0c;对于所有经历过那场鏖战的大学生和指导老师而言&#xff0c;瞬间就能唤起一段刻骨铭心的记忆。它指的并非一所实体大学&#xff0c;而是那场在2021年秋季&#xff0c;…

作者头像 李华
网站建设 2026/8/28 12:44:41

NeSy-RAG:面向可解释问答的神经符号检索增强生成架构

这次我们来看一个比较新的 RAG 方向&#xff1a;NeSy-RAG&#xff0c;全称是 Neuro-Symbolic RAG for Explainable Question Answering 。它要解决的问题很直接&#xff1a;普通 RAG 能告诉你在哪几个文档片段里找到了答案&#xff0c;但很难告诉你这个答案是通过什么逻辑推理…

作者头像 李华
网站建设 2026/8/28 12:44:24

从阶乘约数问题看数论与算法优化:质因数分解与勒让德定理实战

1. 项目概述&#xff1a;从一道竞赛题看算法思维的深度与广度最近在复盘一些经典的算法竞赛题目&#xff0c;蓝桥杯国赛的“阶乘约数”问题又一次引起了我的注意。这道题表面上是一个简单的数论问题&#xff1a;给定一个正整数 n&#xff0c;求 n!&#xff08;n的阶乘&#xff…

作者头像 李华