- 示例工程
- 教程
【免费下载链接】java-design-patterns
Design patterns implemented in Java
数据总线模式(Data Bus Pattern,又称 Event Bus / Message Bus)通过一条集中的消息通道让应用中的各个组件按消息类型收发事件,组件之间互不感知、彻底解耦。本指南以 java-design-patterns 仓库中的data-bus模块为实证蓝本,从模式含义、类结构、完整源码实现、测试验证到实战扩展逐步拆解,读完你将掌握如何在 Java 应用中实现一个支持多对多通信、成员可选择性订阅事件的数据总线。
模式含义:让组件只认“消息类型”
数据总线模式允许一个应用程序的组件之间收发消息/事件,而不需要这些组件相互感知——它们只需要知道所发送/接收的消息/事件的类型即可。它本质上就是常说的 Event-Bus 消息总线模式,属于架构型(Architectural)设计模式,核心标签是解耦(Decoupling)。
仓库中data-bus模块的英文 README 给出了更精确的意图描述:提供一个集中的通信通道,让系统中各个组件无需直接连接即可交换数据,从而促进松耦合,并提升系统的可扩展性与可维护性。
用一句话概括:数据总线把“组件之间的直接通信”替换为“组件与总线之间的通信”,组件只按消息类型决定自己是否接收,彼此之间零依赖。
在 java-design-patterns 的data-bus实现中,DataBus采用单例(Singleton),Member成员注册到总线后即可接收发布到总线上的每一份数据,成员可以自行决定对任何给定消息做出反应或忽略(参见 DataBus.java 与 App.java 的类注释)。该实现是一个同步数据总线:当数据发布到总线时,publish方法在所有成员接收完数据并返回之前不会返回。
现实类比:机场广播系统
可以把一座大型机场想象成数据总线模式的现实映射:航空公司、乘客、行李搬运人员、安保人员都需要沟通与共享信息,但并非每个实体都与其他所有实体直接对话,而是通过一个集中的广播系统(即数据总线)发布航班信息、安全警报等关键更新,各实体只收听与自己相关的消息。这样的设计把通信过程解耦,每个实体只收到自己需要的信息,同时系统可以平滑地接入新实体而不影响既有成员——这正是数据总线模式在 Java 中推动集中式通信与事件处理、提升可扩展性与可维护性的方式。
适用场景
可以在以下场景使用数据总线模式:
- 由组件自己决定接收哪些信息/事件:成员通过
instanceof类型判断选择性消费消息,而非被动接收一切。 - 需要实现多对多的通信:一个总线可以有任意数量的发布者,也可以有任意数量的接收成员,所有成员都会收到同一份数据。
- 希望组件之间不需要感知彼此:发布者不知道谁在订阅,订阅者也不知道谁在发布,双方只依赖消息类型。
英文 README 进一步补充了适用面更广的判断标准:
- 多个组件需要共享数据或事件,但直接耦合不可取;
- 复杂事件驱动系统中信息流动态变化;
- 分布式系统中组件可能部署在不同环境;
- 微服务架构中服务间通信。
模块结构与类图
data-bus模块的代码组织结构如下(完整结构见>/** Members receive events from the Data-Bus. */ public interface Member extends Consumer<DataType> { void accept(DataType event); }
(见 Member.java)
DataType接口要求每个事件都能“知道”自己所在的(或即将发布的)总线,以支持事件在接收处理过程中反向向总线再次发布消息(后续StatusMember的 goodbye 消息会用到这一能力):
/** Events are sent via the Data-Bus. */ public interface DataType { DataBus getDataBus(); void setDataBus(DataBus dataBus); }(见 DataType.java)
AbstractDataType使用 Lombok 的@Getter/@Setter为所有具体事件提供dataBus字段的默认实现,让事件类型可以自由扩展(见 AbstractDataType.java)。
三种事件类型
data包中提供了三个示例DataType实现:
| 事件类 | 承载字段 | 静态工厂方法 | 语义 |
|---|---|---|---|
MessageData | String message(final) | MessageData.of(String message) | 字符串消息事件 |
StartingData | LocalDateTime when(final) | StartingData.of(LocalDateTime when) | 应用启动事件,携带启动时间 |
StoppingData | LocalDateTime when(final) | StoppingData.of(LocalDateTime when) | 应用停止事件,携带停止时间 |
三者均继承AbstractDataType并通过静态工厂方法of(...)构造,源码见 MessageData.java、StartingData.java、StoppingData.java。
总线核心:DataBus 单例的订阅与发布
DataBus是模式的中枢,代码精简但职责清晰(完整实现见 DataBus.java):
public class DataBus { private static final DataBus INSTANCE = new DataBus(); private final Set<Member> listeners = new HashSet<>(); public static DataBus getInstance() { return INSTANCE; } public void subscribe(final Member member) { this.listeners.add(member); } public void unsubscribe(final Member member) { this.listeners.remove(member); } public void publish(final DataType event) { event.setDataBus(this); listeners.forEach(listener -> listener.accept(event)); } }三个方法对应模式的三个基本操作,实现要点如下:
- 单例获取:
INSTANCE在类加载时创建,通过getInstance()全局共享同一个总线实例,保证所有组件接入同一通信中枢。 - 订阅
subscribe(member):将成员加入HashSet监听器集合。使用Set意味着同一成员重复订阅会被去重,且成员间无固定顺序——文档也明确指出“所有成员都会收到同一份数据,但每个成员收到某条数据的顺序是实现细节”。 - 退订
unsubscribe(member):从集合中移除成员,之后该成员不再收到任何事件。 - 发布
publish(event):先把事件自身的dataBus指针指向当前总线(这正是StoppingData能反向发布 goodbye 消息的前提),再遍历全部监听器依次调用accept(event)。这是一个同步、阻塞的发布过程:所有成员处理完毕前publish不会返回。
关键调用链与测试验证
DataBusTest(见 DataBusTest.java)用 Mockito 验证了订阅/退订的核心契约:
@Test void publishedEventIsReceivedBySubscribedMember() { // given final var dataBus = DataBus.getInstance(); dataBus.subscribe(member); // when dataBus.publish(event); // then then(member).should().accept(event); } @Test void publishedEventIsNotReceivedByMemberAfterUnsubscribing() { // given final var dataBus = DataBus.getInstance(); dataBus.subscribe(member); dataBus.unsubscribe(member); // when dataBus.publish(event); // then then(member).should(never()).accept(event); }两条测试用例分别断言:订阅后的成员一定能收到发布的事件;退订后的成员一定收不到后续事件。这正是publish中listeners.forEach(listener -> listener.accept(event))的直接行为证明。
成员实现:按类型选择性接收消息
成员决定接收哪些消息的方式是数据总线模式的精髓——每个成员在自己的accept(DataType data)中用instanceof判断消息类型,只处理自己关心的类型。
普通成员:MessageCollectorMember
普通社区成员只接受MessageData类型的消息,把收到的字符串收集进内部列表(见 MessageCollectorMember.java):
@Slf4j public class MessageCollectorMember implements Member { private final String name; private final List<String> messages = new ArrayList<>(); public MessageCollectorMember(String name) { this.name = name; } @Override public void accept(final DataType data) { if (data instanceof MessageData) { handleEvent((MessageData) data); } } private void handleEvent(MessageData data) { LOGGER.info("{} sees message {}", name, data.getMessage()); messages.add(data.getMessage()); } public List<String> getMessages() { return List.copyOf(messages); } }注意getMessages()返回的是List.copyOf(...)生成的不可变副本,外部无法篡改已收集的消息列表。
状态成员:StatusMember
事件管理员/组织者只接受StartingData与StoppingData两类消息,记录应用的启动/停止时间(见 StatusMember.java):
@Getter @Slf4j @RequiredArgsConstructor public class StatusMember implements Member { private final int id; private LocalDateTime started; private LocalDateTime stopped; @Override public void accept(final DataType data) { if (data instanceof StartingData) { handleEvent((StartingData) data); } else if (data instanceof StoppingData) { handleEvent((StoppingData) data); } } private void handleEvent(StartingData data) { started = data.getWhen(); LOGGER.info("Receiver {} sees application started at {}", id, started); } private void handleEvent(StoppingData data) { stopped = data.getWhen(); LOGGER.info("Receiver {} sees application stopping at {}", id, stopped); LOGGER.info("Receiver {} sending goodbye message", id); data.getDataBus().publish(MessageData.of(String.format("Goodbye cruel world from #%d!", id))); } }StatusMember有一个值得注意的反向发布能力:当收到StoppingData时,它利用事件上的data.getDataBus()拿到当前总线引用,再向总线发布一条MessageData("Goodbye cruel world from #N!")。这生动演示了成员既是订阅者、也可以是发布者,且二者无需知道彼此的存在——这正是多对多通信的体现。
StatusMemberTest(见 StatusMemberTest.java)验证了其类型选择性:传入StartingData会正确记录started时间;传入StoppingData会正确记录stopped时间;传入MessageData时started与stopped均保持为null——说明状态成员确实忽略不属于自己关注类型的消息。
完整运行示例:App 演示与输出解读
App类(见 App.java)串起了整个模式的实际运行流程:
class App { public static void main(String[] args) { final var bus = DataBus.getInstance(); bus.subscribe(new StatusMember(1)); bus.subscribe(new StatusMember(2)); final var foo = new MessageCollectorMember("Foo"); final var bar = new MessageCollectorMember("Bar"); bus.subscribe(foo); bus.publish(StartingData.of(LocalDateTime.now())); bus.publish(MessageData.of("Only Foo should see this")); bus.subscribe(bar); bus.publish(MessageData.of("Foo and Bar should see this")); bus.unsubscribe(foo); bus.publish(MessageData.of("Only Bar should see this")); bus.publish(StoppingData.of(LocalDateTime.now())); } }运行流程可以拆解为五个阶段:
- 订阅:两个
StatusMember(id=1、2)与MessageCollectorMember("Foo")先后注册到总线。 - 发布启动事件:
StartingData广播出去,只有两个StatusMember响应(记录并打印启动时间),Foo 因类型不符而无感知。 - 发布仅 Foo 可见的消息:此时 Bar 尚未订阅,只有 Foo 收到并收集 "Only Foo should see this"。
- Bar 订阅后发布共享消息:"Foo and Bar should see this" 被 Foo 与 Bar 同时收集。
- Foo 退订后发布仅 Bar 消息:Foo 已
unsubscribe,只有 Bar 收到 "Only Bar should see this"。 - 发布停止事件:两个
StatusMember记录停止时间,并各自向总线反向发布一条MessageDatagoodbye 消息——此时仍在订阅中的 Foo、Bar 会收集到这些消息。
当数据总线发布消息时,运行输出大致如下(时间戳因运行时刻而异):
02:33:57.627 [main] INFO com.iluwatar.databus.members.StatusMember - Receiver 2 sees application started at 2022-10-26T02:33:57.613529100 02:33:57.633 [main] INFO com.iluwatar.databus.members.StatusMember - Receiver 1 sees application started at 2022-10-26T02:33:57.613529100如输出所示,MessageCollectorMember只接受MessageData类型,看不到StartingData/StoppingData(它们只对StatusMember可见),从而防止普通成员收到管理员的运行状态通知。这正是“选择性消息处理”的直观体现。
运行方式
data-bus是标准 Maven 模块,与仓库其他模块共用根目录的pom.xml与 Maven Wrapper。可以执行以下命令在模块内运行演示程序与测试(仓库为只读,以下均为本地查看/运行方式):
# 运行演示程序 ./mvnw -pl>赞- 示例工程
- 教程
【免费下载链接】java-design-patterns
Design patterns implemented in Java
相关推荐
深入解析 Java 双缓冲模式(Double Buffer Pattern):以 java-design-patterns 仓库为例
深入解析 Java 双缓冲模式(Double Buffer Pattern):以 java design patterns 仓库为例 双缓冲(Double Bu
示例工程教程Java 数据局部性模式(Data Locality Pattern)深入解析:以 java-design-patterns 源码为实例
Java 数据局部性模式(Data Locality Pattern)深入解析:以 java design patterns 源码为实例 本指南基于 java
示例工程教程Java 领域模型模式(Domain Model Pattern)实战指南:以 java-design-patterns 仓库为例
Java 领域模型模式(Domain Model Pattern)实战指南:以 java design patterns 仓库为例 领域模型模式(Domain
示例工程教程