# EventBus **Repository Path**: vipkwds/eventbus ## Basic Information - **Project Name**: EventBus - **Description**: 高性能、线程安全的 Go 语言事件总线(发布-订阅)系统 - **Primary Language**: Go - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-05-22 - **Last Updated**: 2026-05-22 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # EventBus 事件总线 [![Go Version](https://img.shields.io/badge/Go-1.24.4-blue.svg)](https://github.com/golang/go) [![License](https://img.shields.io/badge/License-MIT-yellow.svg)](LICENSE) > 高性能、线程安全的 Go 语言事件总线(发布-订阅)系统 --- ## 目录 - [特性](#特性) - [快速开始](#快速开始) - [核心概念](#核心概念) - [API 文档](#api-文档) - [使用示例](#使用示例) - [架构设计](#架构设计) - [执行模式](#执行模式) - [错误处理](#错误处理) - [指标监控](#指标监控) - [调用链追踪](#调用链追踪) - [最佳实践](#最佳实践) - [性能优化](#性能优化) - [常见问题](#常见问题) --- ## 特性 ### 核心特性 | 特性 | 描述 | |------|------| | **同步/异步执行** | 支持三种执行模式:同步、异步、异步等待 | | **线程安全** | 全面使用读写锁,支持高并发读操作 | | **泛型结果解析** | 编译期类型安全,避免运行时 panic | | **错误处理器链** | 支持多个错误处理器链式执行 | | **执行指标** | 内置执行次数、成功率、平均耗时监控 | | **调用链追踪** | 支持 TraceID 全链路追踪 | | **处理器去重** | 同一事件类型下相同 ID 的处理器不会重复添加 | | **热迁移** | 支持将旧分发器的处理器迁移到新分发器 | ### 设计亮点 - **零依赖**:核心模块仅依赖 Go 标准库 - **单例模式**:全局单一分发器,简化使用 - **快照拷贝**:发布时拷贝处理器列表,避免锁竞争 - **饿汉初始化**:程序启动时立即初始化,避免懒加载问题 --- ## 快速开始 ### 安装 ```bash go get gitee.com/vipkwds/eventbus ``` ### 最小示例 ```go package main import ( "fmt" "gitee.com/vipkwds/eventbus" ) func main() { // 1. 订阅事件(使用包级便捷函数) eventbus.Subscribe("user.created", eventbus.NewFuncHandler("handler1", func(data interface{}) (interface{}, error) { fmt.Println("收到事件:", data) return "处理成功", nil })) // 2. 发布事件 results, err := eventbus.Publish("user.created", map[string]any{"name": "张三", "age": 25}) // 3. 处理结果 for _, res := range results { fmt.Printf("处理器: %s, 成功: %v, 结果: %v\n", res.HandlerID, res.Success, res.Result) } } ``` 输出: ``` 收到事件: map[name:张三 age:25] 处理器: handler1, 成功: true, 结果: 处理成功 ``` --- ## 核心概念 ### 1. EventHandler(事件处理器) 事件处理的核心接口,所有处理器必须实现此接口: ```go type EventHandler interface { Handle(data interface{}, txMgr ...*utils.TxManager) (interface{}, error) ID() string } ``` **内置处理器实现:** | 类型 | 说明 | |------|------| | `FuncHandler` | 函数适配器,支持返回值 `(interface{}, error)` | | `FuncHandlerWithoutResult` | 兼容旧版无返回值的函数 `(data, txMgr) error` | ### 2. Dispatcher(分发器) 核心调度组件,负责管理订阅和发布: ```go type Dispatcher interface { Subscribe(eventType string, handler EventHandler) Publish(eventType string, data interface{}, txMgr ...*utils.TxManager) ([]ExecutionResult, error) PublishWithMode(eventType string, data interface{}, mode ExecMode) ([]ExecutionResult, error) Unsubscribe(eventType string, handlerID string) bool Has(eventType string) bool Close() } ``` ### 3. Executor(执行器) 负责实际执行已订阅的处理器: ```go type Executor interface { Execute(eventType string, handlers []EventHandler, data interface{}, txMgr ...*utils.TxManager) ([]ExecutionResult, error) } ``` **三种执行器:** | 执行器 | 说明 | |--------|------| | `syncExecutor` | 同步执行,依次调用每个处理器 | | `asyncExecutor` | 异步执行,使用信号量限制并发数 | | `smartExecutor` | 智能调度,根据配置自动选择执行模式 | ### 4. ExecutionResult(执行结果) ```go type ExecutionResult struct { HandlerID string // 处理器唯一标识 Result interface{} // 处理返回的结果 Error error // 执行错误(如果有) Success bool // 是否执行成功 } ``` --- ## API 文档 ### 包级便捷函数 > 位置:`event.go` #### Subscribe - 订阅事件 ```go func Subscribe(eventType string, handler EventHandler) ``` **说明:** 将处理器注册到指定事件类型。 **参数:** | 参数 | 类型 | 说明 | |------|------|------| | eventType | string | 事件类型名称 | | handler | EventHandler | 事件处理器 | **示例:** ```go handler := eventbus.NewFuncHandler("myHandler", func(data interface{}) (interface{}, error) { return nil, nil }) eventbus.Subscribe("order.paid", handler) ``` --- #### SubscribeFunc - 函数订阅(带返回值) ```go func SubscribeFunc(eventType, handlerID string, fn func(data interface{}, txMgr ...*utils.TxManager) (interface{}, error)) ``` **说明:** 直接使用函数作为处理器,简化订阅流程。 **示例:** ```go eventbus.SubscribeFunc("user.registered", "sendWelcomeEmail", func(data interface{}) (interface{}, error) { user := data.(map[string]interface{}) email := user["email"].(string) err := sendEmail(email) return map[string]bool{"emailSent": err == nil}, err }) ``` --- #### SubscribeFuncWithoutResult - 函数订阅(无返回值) ```go func SubscribeFuncWithoutResult(eventType, handlerID string, fn func(data interface{}, txMgr ...*utils.TxManager) error) ``` **说明:** 兼容旧版的无返回值函数订阅。 **示例:** ```go eventbus.SubscribeFuncWithoutResult("order.created", "logOrder", func(data interface{}) error { order := data.(map[string]interface{}) fmt.Printf("新订单: %v\n", order) return nil }) ``` --- #### Publish - 发布事件 ```go func Publish(eventType string, data any, txMgr ...*utils.TxManager) ([]ExecutionResult, error) ``` **说明:** 发布事件到指定类型,所有订阅的处理器将被调用。 **参数:** | 参数 | 类型 | 说明 | |------|------|------| | eventType | string | 事件类型名称 | | data | any | 传递给处理器的数据 | | txMgr | TxManager | 可选,事务管理器 | **返回值:** | 返回值 | 类型 | 说明 | |--------|------|------| | results | []ExecutionResult | 每个处理器的执行结果 | | error | error | 发布过程的错误(通常为 nil) | **示例:** ```go results, err := eventbus.Publish("order.paid", map[string]any{ "orderID": "ORD-12345", "amount": 99.99, "customerID": "CUST-001", }) // 过滤成功/失败结果 successResults := eventbus.GetSuccessfulResults(results) failedResults := eventbus.GetFailedResults(results) ``` --- #### PublishWithMode - 指定模式发布 ```go func PublishWithMode(eventType string, data any, mode ExecMode) ([]ExecutionResult, error) ``` **说明:** 以指定执行模式发布事件。 **mode 可选值:** | 模式 | 值 | 说明 | |------|---|------| | `ModeSync` | 0 | 同步执行(默认) | | `ModeAsync` | 1 | 异步执行(不等待,立即返回空结果) | | `ModeAsyncWait` | 2 | 异步执行但等待所有处理器完成 | **示例:** ```go // 异步发布,不阻塞 results, err := eventbus.PublishWithMode("largeTask", data, ModeAsync) ``` --- #### MustPublish - 发布并忽略错误 ```go func MustPublish(eventType string, data any, txMgr ...*utils.TxManager) ``` **说明:** 发布事件但忽略所有错误和结果,用于 fire-and-forget 场景。 **示例:** ```go // 日志记录等不需要关注结果的场景 eventbus.MustPublish("audit.log", map[string]any{"action": "user.login"}) ``` --- #### Unsubscribe - 取消订阅 ```go func Unsubscribe(eventType string, handlerID string) bool ``` **说明:** 从指定事件类型移除指定的处理器。 **返回值:** `true` 表示成功移除,`false` 表示未找到该处理器。 **示例:** ```go removed := eventbus.Unsubscribe("order.paid", "sendEmailHandler") if removed { fmt.Println("处理器已移除") } ``` --- #### Has - 检查事件是否有订阅者 ```go func Has(eventType string) bool ``` **示例:** ```go if eventbus.Has("order.paid") { fmt.Println("该事件有订阅者") } ``` --- #### Close - 关闭分发器 ```go func Close() ``` **说明:** 关闭分发器,清空所有处理器。通常在程序退出时调用。 --- ### 结果处理函数 #### GetSuccessfulResults - 过滤成功结果 ```go func GetSuccessfulResults(results []ExecutionResult) []ExecutionResult ``` --- #### GetFailedResults - 过滤失败结果 ```go func GetFailedResults(results []ExecutionResult) []ExecutionResult ``` --- #### GetResultByHandlerID - 根据处理器ID获取结果 ```go func GetResultByHandlerID(results []ExecutionResult, handlerID string) (ExecutionResult, bool) ``` --- ### 泛型结果解析 #### ParseEventResult - 泛型结果解析 ```go func ParseEventResult[R any](eventName string, eventResults []ExecutionResult, publishErr error, handlerID ...string) (R, int, error) ``` **泛型约束:** `R` 可以是任意类型(`bool`、`int`、结构体、`map[string]float64` 等) **handlerID 参数:** 可变参数,不传或传空时默认使用 `eventName` 作为匹配 ID。 **错误码说明:** | 错误码 | 常量 | 说明 | |--------|------|------| | 0 | `EventParseCodeSuccess` | 解析成功 | | 1 | `EventParseCodePublish` | 事件发布失败(`publishErr != nil`) | | 2 | `EventParseCodeExec` | 处理器执行失败(`Success = false`) | | 3 | `EventParseCodeType` | 返回结果类型不匹配 | | 4 | `EventParseCodeNoMatch` | 未匹配到指定的 HandlerID | **返回值:** | 返回值 | 说明 | |--------|------| | R | 解析成功=真实业务结果,失败=对应类型的零值 | | int | 错误码,0=成功,其他为失败场景 | | error | 解析成功=nil,失败=精准错误信息 | **完整示例:** ```go // 订阅 eventbus.Subscribe("user.registered", eventbus.NewFuncHandler("user.registered", func(data interface{}) (interface{}, error) { return map[string]interface{}{ "userID": "U12345", "username": "张三", }, nil })) // 发布 results, err := eventbus.Publish("user.registered", map[string]any{"email": "zhangsan@example.com"}) // 泛型解析(自动类型安全) userInfo, code, parseErr := eventbus.ParseEventResult[map[string]interface{}]("user.registered", results, err) if parseErr != nil { switch code { case eventbus.EventParseCodePublish: fmt.Println("发布失败:", parseErr) case eventbus.EventParseCodeExec: fmt.Println("处理器执行失败:", parseErr) case eventbus.EventParseCodeType: fmt.Println("类型不匹配:", parseErr) case eventbus.EventParseCodeNoMatch: fmt.Println("未找到处理器:", parseErr) } } else { fmt.Printf("用户注册成功: %+v\n", userInfo) } ``` --- ## 使用示例 ### 示例 1:基本发布-订阅 ```go package main import ( "fmt" "gitee.com/vipkwds/eventbus" ) func main() { // 定义处理器 handler1 := eventbus.NewFuncHandler("handler1", func(data interface{}) (interface{}, error) { order := data.(map[string]interface{}) fmt.Printf("[Handler1] 处理订单: %s, 金额: %.2f\n", order["id"], order["amount"]) return "handler1 done", nil }) handler2 := eventbus.NewFuncHandler("handler2", func(data interface{}) (interface{}, error) { order := data.(map[string]interface{}) fmt.Printf("[Handler2] 发送邮件通知: %s\n", order["email"]) return "email sent", nil }) // 订阅 eventbus.Subscribe("order.created", handler1) eventbus.Subscribe("order.created", handler2) // 发布 orderData := map[string]interface{}{ "id": "ORD-001", "amount": 199.99, "email": "customer@example.com", } results, _ := eventbus.Publish("order.created", orderData) fmt.Println("\n--- 执行结果 ---") for _, res := range results { fmt.Printf("处理器: %s, 成功: %v\n", res.HandlerID, res.Success) } } ``` 输出: ``` [Handler1] 处理订单: ORD-001, 金额: 199.99 [Handler2] 发送邮件通知: customer@example.com --- 执行结果 --- 处理器: handler1, 成功: true 处理器: handler2, 成功: true ``` --- ### 示例 2:异步执行 ```go func main() { // 订阅一个耗时操作 eventbus.Subscribe("report.generate", eventbus.NewFuncHandler("generator", func(data interface{}) (interface{}, error) { time.Sleep(2 * time.Second) // 模拟耗时操作 return "report ready", nil })) fmt.Println("开始异步发布...") // 异步发布,不等待完成 results, _ := eventbus.PublishWithMode("report.generate", nil, eventbus.ModeAsync) fmt.Println("立即返回,继续执行其他任务...") fmt.Printf("异步结果(空): %v\n", results) // 返回空切片 time.Sleep(3 * time.Second) // 模拟等待 fmt.Println("任务完成") } ``` --- ### 示例 3:异步等待 ```go func main() { // 订阅多个处理器 for i := 1; i <= 3; i++ { idx := i eventbus.Subscribe("task.parallel", eventbus.NewFuncHandler(fmt.Sprintf("task%d", idx), func(data interface{}) (interface{}, error) { time.Sleep(1 * time.Second) return fmt.Sprintf("task%d done", idx), nil })) } fmt.Println("开始并行执行...") // 异步等待模式 results, _ := eventbus.PublishWithMode("task.parallel", nil, eventbus.ModeAsyncWait) fmt.Println("所有任务完成,结果:") for _, res := range results { fmt.Printf(" %s: %v\n", res.HandlerID, res.Result) } } ``` --- ### 示例 4:使用泛型解析结果 ```go package main import ( "fmt" "gitee.com/vipkwds/eventbus" ) type UserResponse struct { UserID string Username string Email string } func main() { // 订阅 - 返回结构体 eventbus.Subscribe("user.get", eventbus.NewFuncHandler("user.get", func(data interface{}) (interface{}, error) { return UserResponse{ UserID: "U10086", Username: "张三", Email: "zhangsan@example.com", }, nil })) // 发布 results, err := eventbus.Publish("user.get", map[string]any{"userID": "U10086"}) // 使用泛型精确解析 var user UserResponse user, code, parseErr := eventbus.ParseEventResult[UserResponse]("user.get", results, err) if parseErr != nil { fmt.Printf("解析失败: %v (code: %d)\n", parseErr, code) return } fmt.Printf("用户信息: ID=%s, Name=%s, Email=%s\n", user.UserID, user.Username, user.Email) } ``` --- ### 示例 5:自定义错误处理器 ```go package main import ( "fmt" "gitee.com/vipkwds/eventbus" ) // 自定义错误处理器 type CustomErrorHandler struct{} func (h *CustomErrorHandler) HandleError(eventType, handlerID string, err error) { fmt.Printf("[自定义错误] 事件: %s, 处理器: %s, 错误: %v\n", eventType, handlerID, err) } // 触发错误的处理器 func main() { // 创建带自定义错误处理器的分发器 customHandler := &CustomErrorHandler{} dispatcher := eventbus.NewDispatcher(eventbus.DefaultExecutorConfig, customHandler) // 使用自定义分发器 eventbus.SetDispatcher(dispatcher) // 订阅 eventbus.Subscribe("error.test", eventbus.NewFuncHandler("errorHandler", func(data interface{}) (interface{}, error) { return nil, fmt.Errorf("故意返回的错误") })) // 发布(处理器返回错误) results, _ := eventbus.Publish("error.test", nil) // 查看结果中的错误 for _, res := range results { if res.Error != nil { fmt.Printf("执行错误: %v\n", res.Error) } } } ``` --- ### 示例 6:指标监控 ```go func main() { // 执行多次 eventbus.Subscribe("metric.test", eventbus.NewFuncHandler("handler1", func(data interface{}) (interface{}, error) { return "ok", nil })) for i := 0; i < 10; i++ { eventbus.Publish("metric.test", nil) } // 获取指标 metrics := eventbus.GetMetrics("metric.test") fmt.Printf("总执行次数: %d\n", metrics.TotalExecutions) fmt.Printf("成功次数: %d\n", metrics.SuccessCount) fmt.Printf("失败次数: %d\n", metrics.FailureCount) fmt.Printf("成功率: %.2f%%\n", metrics.GetSuccessRate()*100) fmt.Printf("平均耗时: %.2fms\n", metrics.GetAvgDurationMs()) // 获取所有事件的指标 allMetrics := eventbus.GetAllMetrics() for eventType, m := range allMetrics { fmt.Printf("事件 %s: 执行%d次\n", eventType, m.TotalExecutions) } } ``` --- ### 示例 7:调用链追踪 ```go func main() { // 创建追踪上下文 trace := eventbus.NewTrace("order.process") // 创建可追踪的处理器 traceableHandler := eventbus.NewTraceableHandler( eventbus.NewFuncHandler("orderHandler", func(data interface{}) (interface{}, error) { return "processed", nil }), trace, ) // 订阅并发布 eventbus.Subscribe("order.process", traceableHandler) eventbus.Publish("order.process", map[string]any{"orderID": "123"}) fmt.Printf("追踪ID: %s, 耗时: %v\n", trace.TraceID, trace.Duration()) } ``` --- ### 示例 8:处理器热迁移 ```go func main() { // 在旧分发器上订阅 eventbus.Subscribe("legacy.event", eventbus.NewFuncHandler("oldHandler", func(data interface{}) (interface{}, error) { return "old", nil })) // 创建新分发器 newDispatcher := eventbus.NewDispatcher(eventbus.DefaultExecutorConfig) // 迁移处理器 eventbus.MigrateHandlers(newDispatcher) // 设置使用新分发器 eventbus.SetDispatcher(newDispatcher) // 验证迁移成功 results, _ := eventbus.Publish("legacy.event", nil) fmt.Printf("迁移后结果: %+v\n", results) } ``` --- ## 架构设计 ### 整体架构图 ``` ┌─────────────────────────────────────────────────────────────────┐ │ eventbus │ ├─────────────────────────────────────────────────────────────────┤ │ │ │ ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │ │ │ Subscribe │────▶│ Dispatcher │────▶│ Executor │ │ │ │ (event) │ │ (registry) │ │ │ │ │ └─────────────┘ └──────────────┘ │ ┌───────────┐ │ │ │ │ │ │ │ syncExec │ │ │ │ │ │ │ ├───────────┤ │ │ │ ┌─────────────┐ ┌──────────────┐ │ │ asyncExec │ │ │ │ │ Publish │────▶│ Handler │────▶│ ├───────────┤ │ │ │ │ (event) │ │ Registry │ │ │smartExec │ │ │ │ └─────────────┘ └──────────────┘ │ └───────────┘ │ │ │ └─────────────────┘ │ │ │ │ ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │ │ │ Metrics │◀────│ Collector │◀────│ Recorder │ │ │ │ (metrics) │ │ │ │ │ │ │ └─────────────┘ └──────────────┘ └─────────────────┘ │ │ │ │ ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │ │ │ Trace │◀────│ TraceContext │◀────│ TraceableHandler│ │ │ │ (trace) │ │ │ │ │ │ │ └─────────────┘ └──────────────┘ └─────────────────┘ │ │ │ │ ┌─────────────────────────────────────────────────────────────┐│ │ │ ErrorHandler ││ │ │ ┌─────────────┐ ┌─────────────┐ ┌────────────────────┐ ││ │ │ │LogErrorHandler│ │ErrorHandlerChain│ │CustomErrorHandler │ ││ │ │ └─────────────┘ └─────────────┘ └────────────────────┘ ││ │ └─────────────────────────────────────────────────────────────┘│ └─────────────────────────────────────────────────────────────────┘ ``` ### 模块关系 ``` 用户代码 │ ▼ ┌──────────────────────────────────────────────────────────────────┐ │ 包级便捷函数 (event.go) │ │ Subscribe() / Publish() / Unsubscribe() / GetSuccessfulResults()│ └──────────────────────────────────────────────────────────────────┘ │ │ ▼ ▼ ┌──────────────────┐ ┌───────────────────────────────────────────┐ │ registry.go │ │ dispatcher.go │ │ │ │ │ │ 全局单例管理 │ │ ┌─────────────────────────────────────┐ │ │ - globalDispatcher│ │ │ defaultDispatcher │ │ │ - GetDispatcher()│ │ │ │ │ │ - SetDispatcher()│ │ │ handlers: map[string][]EventHandler │ │ │ - MigrateHandlers│ │ │ executor: Executor │ │ └──────────────────┘ │ │ errorHandler: ErrorHandler │ │ │ │ mu: sync.RWMutex │ │ │ └─────────────────────────────────────┘ │ └───────────────────────────────────────────┘ │ ▼ ┌───────────────────────┐ │ executor.go │ │ │ │ ┌─────────────────┐ │ │ │ syncExecutor │ │ │ └─────────────────┘ │ │ ┌─────────────────┐ │ │ │ asyncExecutor │ │ │ └─────────────────┘ │ │ ┌─────────────────┐ │ │ │ smartExecutor │ │ │ └─────────────────┘ │ └───────────────────────┘ │ ▼ ┌───────────────────────┐ │ handler.go │ │ │ │ EventHandler 接口 │ │ FuncHandler │ │ FuncHandlerWithout │ └───────────────────────┘ ``` ### 线程安全设计 #### 1. 读写锁分离 ```go type defaultDispatcher struct { mu sync.RWMutex // 读写锁 handlers map[string][]EventHandler // ... } ``` - **读操作**(`Subscribe`、`Publish`、`Unsubscribe`、`Has`):使用 `RLock()` - **写操作**(`Close`):使用 `Lock()` #### 2. 快照拷贝模式 ```go func (d *defaultDispatcher) Publish(...) { d.mu.RLock() handlers := make([]EventHandler, len(d.handlers[eventType])) copy(handlers, d.handlers[eventType]) d.mu.RUnlock() // 立即释放读锁 // 执行时使用拷贝的切片 executor.Execute(eventType, handlers, data) } ``` **优势:** - 写者(Subscribe/Unsubscribe)不会被读者(Publish)阻塞 - 执行时间长也不会阻塞其他发布操作 #### 3. 单例初始化双重检查 ```go func GetDispatcher() Dispatcher { mu.Lock() defer mu.Unlock() if globalDispatcher == nil { globalDispatcher = NewDispatcher(DefaultExecutorConfig) } return globalDispatcher } ``` #### 4. 原子操作 指标记录使用 `sync/atomic` 保证线程安全: ```go atomic.AddInt64(&m.TotalExecutions, 1) atomic.AddInt64(&m.SuccessCount, 1) atomic.LoadInt64(&m.TotalExecutions) ``` --- ## 执行模式 ### 模式对比 | 模式 | 值 | 行为 | 返回结果 | 适用场景 | |------|---|------|----------|----------| | `ModeSync` | 0 | 同步依次执行所有处理器 | 完整结果集 | 需要等待所有处理完成的业务 | | `ModeAsync` | 1 | 异步启动所有处理器,不等待 | 空切片 `[]` | 日志、通知等不需要结果的操作 | | `ModeAsyncWait` | 2 | 异步执行,等待所有完成 | 完整结果集 | 耗时操作希望并发执行但最终要结果 | ### 执行流程 #### ModeSync 流程 ``` Publish("order.created", data) │ ▼ ┌─────────────────┐ │ 获取读锁 │ │ 拷贝处理器列表 │ └────────┬────────┘ │ ▼ ┌─────────────────┐ │ 释放读锁 │ │ 同步遍历执行 │ │ handler1.Handle() │ handler2.Handle() │ handler3.Handle() └────────┬────────┘ │ ▼ 返回结果集 ``` #### ModeAsync 流程 ``` PublishWithMode("event", data, ModeAsync) │ ▼ ┌─────────────────┐ │ 获取读锁 │ │ 拷贝处理器列表 │ └────────┬────────┘ │ ▼ ┌─────────────────┐ │ 释放读锁 │ │ 启动 goroutine │ │ ┌────────────┐ │ │ │ go h1.Handle│ │ │ ├────────────┤ │ │ │ go h2.Handle│ │ │ ├────────────┤ │ │ │ go h3.Handle│ │ │ └────────────┘ │ └────────┬────────┘ │ ▼ 立即返回 [] ``` #### ModeAsyncWait 流程 ``` PublishWithMode("event", data, ModeAsyncWait) │ ▼ ┌─────────────────┐ │ 获取读锁 │ │ 拷贝处理器列表 │ └────────┬────────┘ │ ▼ ┌─────────────────┐ │ 释放读锁 │ │ 并发执行 │ │ ┌─ semaphore ─┐│ │ │ (最多N个) ││ │ └────────────┘│ │ WaitGroup等待 │ │ timeout控制 │ └────────┬────────┘ │ ▼ 收集所有结果 返回结果集 ``` ### 超时控制 `asyncExecutor` 支持超时配置: ```go config := eventbus.ExecutorConfig{ Mode: eventbus.ModeAsyncWait, MaxWorkers: 50, Timeout: 5000, // 毫秒,5秒超时 } dispatcher := eventbus.NewDispatcher(config) ``` --- ## 错误处理 ### 错误类型 | 类型 | 文件 | 说明 | |------|------|------| | `EventError` | errors.go | 可扩展的错误类型,包含事件类型和处理器ID | | `MultiError` | errors.go | 多个错误的聚合,支持 `Add()` 和 `IsEmpty()` | ### EventError 用法 ```go err := &EventError{ EventType: "order.created", HandlerID: "emailHandler", Err: fmt.Errorf("smtp connection failed"), } fmt.Println(err.Error()) // 输出: 事件[order.created]处理器[emailHandler]错误: smtp connection failed ``` ### MultiError 用法 ```go multiErr := &MultiError{} multiErr.Add(err1) multiErr.Add(err2) if !multiErr.IsEmpty() { fmt.Println(multiErr.Error()) // 输出: 有2个错误发生 } ``` ### 错误处理器 #### ErrorHandler 接口 ```go type ErrorHandler interface { HandleError(eventType string, handlerID string, err error) } ``` #### 内置日志错误处理器 ```go handler := eventbus.NewLogErrorHandler() dispatcher := eventbus.NewDispatcher(config, handler) ``` #### 错误处理器链 ```go // 链式处理多个错误 chain := eventbus.NewErrorHandlerChain( &LogHandler{}, &SlackNotifier{}, &MetricsRecorder{}, ) dispatcher := eventbus.NewDispatcher(config, chain) ``` #### 自定义错误处理器 ```go type MyErrorHandler struct{} func (h *MyErrorHandler) HandleError(eventType, handlerID string, err error) { // 发送到监控系统 sentry.CaptureException(err) // 发送告警 sendAlert(eventType, handlerID, err) } dispatcher := eventbus.NewDispatcher(config, &MyErrorHandler{}) ``` --- ## 指标监控 ### 内置指标 | 指标 | 类型 | 说明 | |------|------|------| | `TotalExecutions` | int64 | 总执行次数 | | `SuccessCount` | int64 | 成功次数 | | `FailureCount` | int64 | 失败次数 | | `TotalDurationNs` | int64 | 总耗时(纳秒) | | `LastExecutionTime` | int64 | 最后执行时间(纳秒时间戳) | ### 计算方法 ```go // 获取平均执行时长(毫秒) avgMs := metrics.GetAvgDurationMs() // 获取成功率 rate := metrics.GetSuccessRate() // 0.0 ~ 1.0 ``` ### 使用示例 ```go func monitorEventHealth() { for { time.Sleep(1 * time.Minute) allMetrics := eventbus.GetAllMetrics() for eventType, m := range allMetrics { if m.GetSuccessRate() < 0.95 { fmt.Printf("[警告] 事件 %s 成功率过低: %.2f%%\n", eventType, m.GetSuccessRate()*100) } } } } ``` --- ## 调用链追踪 ### TraceContext ```go type TraceContext struct { TraceID string // 唯一追踪ID EventType string // 事件类型 StartTime time.Time // 开始时间 } ``` ### 使用方式 ```go func main() { // 创建追踪 trace := eventbus.NewTrace("payment.process") // 包装处理器 traceableHandler := eventbus.NewTraceableHandler( myHandler, trace, ) // 订阅 eventbus.Subscribe("payment.process", traceableHandler) // 发布 eventbus.Publish("payment.process", paymentData) // 输出追踪信息 fmt.Printf("TraceID: %s, Duration: %v\n", trace.TraceID, trace.Duration()) } ``` ### 追踪日志 ```go trace := eventbus.NewTrace("order.create") // 带追踪ID的日志 trace.Log("开始处理订单: %s", orderID) trace.LogError("处理失败: %v", err) ``` --- ## 最佳实践 ### 1. 处理器 ID 命名 ```go // ✅ 推荐:使用有意义的ID eventbus.Subscribe("order.created", NewFuncHandler("inventory.deduct", ...)) eventbus.Subscribe("order.created", NewFuncHandler("email.send", ...)) // ❌ 避免:使用无意义或重复的ID eventbus.Subscribe("order.created", NewFuncHandler("h1", ...)) eventbus.Subscribe("order.created", NewFuncHandler("h1", ...)) // 会去重! ``` ### 2. 错误处理 ```go // ✅ 推荐:检查错误码进行精细化处理 user, code, err := ParseEventResult[User](...) switch code { case EventParseCodeSuccess: // 成功处理 case EventParseCodeExec: // 业务错误处理 case EventParseCodeNoMatch: // 缺少处理器处理 } // ❌ 避免:只检查 err != nil if err != nil { return err } ``` ### 3. 类型断言 ```go // ✅ 推荐:使用泛型自动类型校验 result, _, _ := ParseEventResult[map[string]int](...) // 编译期确保类型安全 // ❌ 避免:运行时断言(可能panic) result := res.Result.(map[string]int) // 危险! ``` ### 4. 事务管理 ```go // ✅ 推荐:同步模式传递事务管理器 results, err := eventbus.Publish("order.create", data, txMgr) // ❌ 注意:异步模式不支持事务 eventbus.PublishWithMode("order.create", data, ModeAsync) // txMgr 在异步模式下会被忽略 ``` ### 5. 资源清理 ```go func main() { defer eventbus.Close() // 程序退出时清理 // ... 业务代码 } ``` ### 6. 订阅位置 ```go // ✅ 推荐:在 init() 或 main() 早期订阅 func init() { eventbus.Subscribe("app.start", handler) } // ❌ 避免:在业务逻辑中动态订阅 func processOrder(order Order) { eventbus.Subscribe("order.processed", handler) // 不推荐 } ``` --- ## 性能优化 ### 1. 读写锁优化 发布操作使用读锁,不会阻塞其他发布操作: ```go d.mu.RLock() // 读锁 handlers = copy(...) // 拷贝 d.mu.RUnlock() // 立即释放 executor.Execute(...) // 执行(无锁) ``` ### 2. 信号量限制并发 异步执行使用信号量防止 goroutine 爆炸: ```go semaphore := make(chan struct{}, e.config.MaxWorkers) // 最多100并发 ``` ### 3. 预分配切片 ```go results := make([]ExecutionResult, 0, len(handlers)) // 预估容量 ``` ### 4. 高频事件优化 如果某事件处理非常频繁,考虑: ```go // 使用异步模式 eventbus.PublishWithMode("high.freq", data, ModeAsync) // 或使用消息队列(如Kafka)处理极高频场景 ``` ### 5. 避免大对象传递 ```go // ❌ 避免:传递大对象 eventbus.Publish("large.data", hugeSlice) // 拷贝整个切片 // ✅ 推荐:传递指针或使用通道 eventbus.Publish("ref.data", &hugeStruct) ``` --- ## 常见问题 ### Q: 为什么 Publish 返回空切片? **A:** 检查以下情况: 1. 事件类型没有订阅者 2. 事件类型拼写错误 3. 使用了 `ModeAsync` 模式(该模式返回空切片) ```go // 调试 if !eventbus.Has("event.name") { fmt.Println("该事件没有订阅者") } ``` ### Q: 处理器为什么没有执行? **A:** 常见原因: 1. 处理器 ID 重复,第二个订阅被忽略 2. 分发器被关闭(`Close()` 后不能再发布) 3. 使用了 `SetDispatcher()` 切换了分发器 ### Q: 异步模式为什么返回空结果? **A:** 这是设计如此。`ModeAsync` 是 fire-and-forget 模式,不等待执行完成。如果需要异步执行但等待结果,使用 `ModeAsyncWait`。 ### Q: 如何调试发布-订阅流程? ```go // 1. 启用追踪 trace := eventbus.NewTrace("debug.event") traceable := eventbus.NewTraceableHandler(handler, trace) eventbus.Subscribe("debug.event", traceable) // 2. 查看指标 fmt.Printf("Metrics: %+v\n", eventbus.GetMetrics("debug.event")) // 3. 检查订阅者数量 fmt.Printf("Has subscribers: %v\n", eventbus.Has("debug.event")) ``` ### Q: 是否支持集群部署? **A:** 当前版本是单节点内存实现。如需集群部署,需要: 1. 外部消息队列(Kafka、RabbitMQ) 2. 或者使用分布式事件总线实现 --- ## 更新日志格式 ```markdown ## [版本号] - 日期 ### 新增 - 新功能描述 ### 优化 - 性能改进 ### 修复 - Bug 修复 ### 破坏性变更 - 不兼容变更说明 ``` --- ## 许可证 本项目基于 MIT 许可证开源。 --- ## 贡献 欢迎提交 Issue 和 Pull Request!