事件总线
一、核心功能
事件总线是 Furion.Pure 框架提供的发布-订阅模式实现,用于解耦应用中的各个模块,实现事件驱动架构。
1.1 核心价值
- 模块解耦:发布者和订阅者互不依赖
- 异步处理:支持异步事件处理,提高系统响应速度
- 多订阅者:同一事件可以有多个订阅者
- 模糊匹配:支持事件名称模糊匹配
- 可扩展性:支持自定义事件源存储和发布者
二、基本概念
2.1 事件源 (EventSource)
事件源是事件的载体,包含事件名称、数据和元信息:
publicinterfaceIEventSource{/// <summary>/// 事件 ID/// </summary>stringEventId{get;}/// <summary>/// 事件名称/// </summary>stringEventName{get;}/// <summary>/// 事件数据/// </summary>objectPayload{get;}/// <summary>/// 事件创建时间/// </summary>DateTimeCreatedTime{get;}/// <summary>/// 是否异步执行/// </summary>boolIsAsync{get;}}2.2 事件发布者 (EventPublisher)
事件发布者负责将事件发送到事件总线:
publicinterfaceIEventPublisher{/// <summary>/// 发布事件/// </summary>/// <param name="eventSource">事件源</param>TaskPublishAsync(IEventSourceeventSource);}2.3 事件订阅者 (EventSubscriber)
事件订阅者负责处理订阅的事件:
publicinterfaceIEventSubscriber{/// <summary>/// 订阅事件/// </summary>/// <param name="eventName">事件名称</param>/// <param name="handler">事件处理器</param>voidSubscribe(stringeventName,Func<EventHandlerExecutingContext,Task>handler);/// <summary>/// 取消订阅/// </summary>/// <param name="eventName">事件名称</param>voidUnsubscribe(stringeventName);}三、实现流程
3.1 服务注册
在Startup.cs中调用:
services.AddEventBus();注册逻辑:
| 步骤 | 操作 | 说明 |
|---|---|---|
| 1 | 注册事件源存储 | 添加ChannelEventSourceStorer(内存通道) |
| 2 | 注册事件发布者 | 添加ChannelEventPublisher |
| 3 | 注册事件总线工厂 | 添加EventBusFactory |
| 4 | 注册后台服务 | 添加EventBusHostedService |
3.2 事件发布流程
┌─────────────────────────────────────────────────────────────┐ │ 事件发布阶段 │ ├─────────────────────────────────────────────────────────────┤ │ 1. 创建事件源 │ │ └── new ChannelEventSource("UserCreated", data) │ ├─────────────────────────────────────────────────────────────┤ │ 2. 调用 IEventPublisher.PublishAsync() │ │ └── 将事件源存储到 ChannelEventSourceStorer │ ├─────────────────────────────────────────────────────────────┤ │ 3. EventBusHostedService 监听通道 │ │ └── 发现新事件源并触发处理 │ └─────────────────────────────────────────────────────────────┘3.3 事件处理流程
┌─────────────────────────────────────────────────────────────┐ │ 事件处理阶段 │ ├─────────────────────────────────────────────────────────────┤ │ 1. 匹配事件订阅者 │ │ └── 根据事件名称查找订阅者 │ ├─────────────────────────────────────────────────────────────┤ │ 2. 创建事件处理上下文 │ │ └── EventHandlerExecutingContext │ ├─────────────────────────────────────────────────────────────┤ │ 3. 执行事件处理器 │ │ └── 调用订阅者注册的处理方法 │ ├─────────────────────────────────────────────────────────────┤ │ 4. 记录处理结果 │ │ └── EventHandlerExecutedContext │ └─────────────────────────────────────────────────────────────┘四、使用示例
4.1 定义事件源
publicclassUserCreatedEvent:ChannelEventSource{publicUserCreatedEvent(UserInfouser):base("UserCreated",user){}}4.2 发布事件
publicclassUserService{privatereadonlyIEventPublisher_eventPublisher;publicUserService(IEventPublishereventPublisher){_eventPublisher=eventPublisher;}publicasyncTaskCreateUser(UserInfouser){// 创建用户逻辑await_eventPublisher.PublishAsync(newUserCreatedEvent(user));}}4.3 订阅事件(特性方式)
使用[EventSubscribe]特性订阅事件:
[EventSubscribe("UserCreated")]publicclassUserCreatedHandler{publicasyncTaskHandle(EventHandlerExecutingContextcontext){varuser=context.PayloadasUserInfo;// 处理用户创建事件}}4.4 订阅事件(代码方式)
使用IEventSubscriber订阅事件:
publicclassEventController{privatereadonlyIEventSubscriber_eventSubscriber;publicEventController(IEventSubscribereventSubscriber){_eventSubscriber=eventSubscriber;}publicvoidSubscribe(){_eventSubscriber.Subscribe("UserCreated",asynccontext=>{varuser=context.PayloadasUserInfo;// 处理用户创建事件});}}五、配置选项
EventBusOptionsBuilder提供了丰富的配置项:
| 配置项 | 默认值 | 说明 |
|---|---|---|
ChannelCapacity | 10000 | 通道容量 |
UseUtcTimestamp | false | 是否使用 UTC 时间 |
FuzzyMatch | false | 是否启用模糊匹配 |
GCCollect | false | 是否启用垃圾回收 |
LogEnabled | true | 是否启用日志 |
5.1 配置示例
services.AddEventBus(options=>{options.ChannelCapacity=10000;options.FuzzyMatch=true;options.LogEnabled=true;});六、高级特性
6.1 模糊匹配
启用模糊匹配后,可以使用通配符订阅事件:
[EventSubscribe("User.*")]publicclassUserEventHandler{publicasyncTaskHandle(EventHandlerExecutingContextcontext){// 处理所有 User 开头的事件}}6.2 异步执行
事件处理器默认异步执行,可以通过IsAsync属性控制:
publicclassUserCreatedEvent:ChannelEventSource{publicUserCreatedEvent(UserInfouser):base("UserCreated",user,isAsync:true){}}6.3 事件监听
实现IEventHandlerMonitor接口监听事件处理:
publicclassEventMonitor:IEventHandlerMonitor{publicvoidOnExecuting(EventHandlerExecutingContextcontext){// 事件处理开始}publicvoidOnExecuted(EventHandlerExecutedContextcontext){// 事件处理完成}}6.4 失败策略
实现IEventFallbackPolicy接口自定义失败处理策略:
publicclassRetryFallbackPolicy:IEventFallbackPolicy{publicasyncTaskHandleAsync(EventHandlerExecutingContextcontext,Exceptionexception){// 重试或其他失败处理逻辑}}6.5 自定义事件源存储
实现IEventSourceStorer接口自定义事件源存储:
publicclassRedisEventSourceStorer:IEventSourceStorer{publicValueTaskWriteAsync(IEventSourceeventSource){// 写入 Redis}publicIAsyncEnumerable<IEventSource>ReadAllAsync(){// 从 Redis 读取}}6.6 消息中心
使用MessageCenter简化事件发布:
// 发布事件awaitMessageCenter.PublishAsync("UserCreated",user);// 订阅事件MessageCenter.Subscribe("UserCreated",async(payload)=>{varuser=payloadasUserInfo;// 处理事件});七、核心文件
| 文件 | 说明 |
|---|---|
EventBusServiceCollectionExtensions.cs | 事件总线服务扩展方法 |
IEventSource.cs | 事件源接口 |
ChannelEventSource.cs | 内存通道事件源 |
IEventPublisher.cs | 事件发布者接口 |
ChannelEventPublisher.cs | 内存通道事件发布者 |
IEventSubscriber.cs | 事件订阅者接口 |
IEventSourceStorer.cs | 事件源存储接口 |
ChannelEventSourceStorer.cs | 内存通道事件源存储 |
EventSubscribeAttribute.cs | 事件订阅特性 |
EventBusHostedService.cs | 事件总线后台服务 |
MessageCenter.cs | 消息中心 |
八、总结
事件总线通过发布-订阅模式实现了模块间的解耦,核心设计思想:
- 发布-订阅模式:发布者和订阅者互不依赖,通过事件总线通信
- 异步处理:默认异步执行事件处理器,提高系统响应速度
- 多订阅者支持:同一事件可以有多个订阅者,实现广播效果
- 模糊匹配:支持事件名称模糊匹配,灵活订阅相关事件
- 高度可扩展:支持自定义事件源存储、发布者和失败策略
这种设计使得应用中的各个模块可以独立开发和测试,提高了系统的可维护性和扩展性。