go中可用channel实现轻量级事件总线,但仅适用于单进程、无持久化、不需顺序与重试的场景;事件溯源须存含id、aggregateid、version等字段的不可变操作日志,并用postgresql/sqlite存储以支持重放与快照。

Go 里没有现成的“事件总线 + 事件溯源”开箱即用框架,但可以用标准库和少量结构体自己搭出来——关键不是套概念,而是明确哪些状态必须持久化、哪些通知可以丢、以及事件消费失败时要不要重试。
用 chan 做轻量级异步事件总线够不够用?
够,但仅限于单进程、无持久化、不关心投递顺序和失败重试的场景。比如服务内部模块间松耦合通知(用户注册成功后发邮件、写日志)。
实操建议:
- 用
chan interface{}或带类型的 channel(如chan *UserCreatedEvent)做广播通道,配合select+default实现非阻塞发送 - 每个订阅者起独立 goroutine 消费,避免一个卡住拖垮全部;别直接在发布端同步调用 handler
- 别把 channel 当队列用:它不保证顺序(多生产者时)、不支持回溯、满载会 panic 或阻塞——需要背压控制就得自己加 buffer 或用
select判断 - 如果要支持多个 topic,别用 map[topic]chan —— channel 本身不能 close 后复用,应封装成
type EventBus struct { mu sync.RWMutex; topics map[string]map[func(interface{})]struct{} }这类结构
事件溯源的核心不是存 JSON,而是存“可重放的操作日志”
事件溯源(Event Sourcing)本质是把状态变更记录为不可变事件序列,再通过重放重建当前状态。Go 里最容易踩的坑是:把事件当普通日志存,却没留出 Version、AggregateID、EventType 等必要字段。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
实操建议:
- 每个事件结构体必须含:
ID string(全局唯一,推荐 ULID)、AggregateID string(聚合根 ID)、Version uint64(乐观并发控制用)、Timestamp time.Time、EventType string(如"user_created") - 存储层选 PostgreSQL(带
jsonb和事务)或 SQLite(嵌入式场景),避免用纯文件追加——你得支持按AggregateID查询历史、支持事务内原子写事件+更新快照 - 别在事件里存指针或闭包;所有字段必须可序列化(
json.Marshal能过);敏感字段(如密码)绝不能进事件 - 快照(Snapshot)不是可选优化项:当某个聚合重放超过 100 条事件时,性能会明显下降,应在
Version % 100 == 0时保存一次当前状态的序列化值
github.com/ThreeDotsLabs/watermill 能直接拿来跑事件溯源吗?
不能。Watermill 是消息中间件适配层(Kafka/RabbitMQ/Redis),它解决的是“跨服务事件分发”,不是“单服务内事件持久化与状态重建”。强行套用会导致两个问题:事件丢失(MQ 不保证 at-least-once 对单个消费者)、版本错乱(不同实例重放进度不一致)。
实操建议:
- 如果系统已用 Kafka,可用它做「事件外发」(如通知其他服务),但「事件溯源主存储」仍需走本地 DB;用 Watermill 的
Publisher在事件写 DB 成功后再发一份到 Kafka - 别让消费者直接从 Kafka 重放重建状态——网络延迟、分区重平衡会让
Offset不可靠;重建逻辑必须只依赖本地事件表 - Watermill 的
Handler默认不提供幂等性保障,重复消费时若没校验AggregateID + Version,可能造成状态覆盖错误
真正难的不是发事件或存事件,而是定义清楚“哪个操作算一个原子事件”——比如转账操作,到底是 TransferStarted、BalanceDeducted、BalanceAdded 分三步发,还是只发一个 TransferCompleted?这决定了你的业务一致性边界在哪,也决定了补偿逻辑写不写得出来。










