Go事件溯源:Event Sourcing实现

发布时间:2026/8/19 21:07:51
Go事件溯源:Event Sourcing实现 Go事件溯源:Event Sourcing实现摘要: 本篇讲解Go语言Event Sourcing实现设计Event Store持久化事件流用聚合根快照加速状态重建事件回放重建聚合状态幂等处理重复事件分享事件版本迁移导致老数据无法回放的踩坑经验对比Event Sourcing、CRUD、快照三种状态管理方案。开篇故事去年支付系统出过一次事故。一个账户余额从1000变成了-50但没人知道怎么变的。查数据库只有最终状态操作日志只记录余额更新没有变更过程。我们对账对了一整夜怀疑是并发扣款没加锁。这事让我下定决心引入事件溯源。把每次状态变更都存成不可变的事件账户余额由事件回放计算。出了问题把事件流拉出来一看就清楚哪一笔扣款在什么时候发生余额怎么一步步变化的。这篇把Event Store、聚合快照、事件回放、版本迁移都写清楚。一、Event Store设计与事件流Event Sourcing的核心是状态由事件推导。传统方式存最终状态Event Sourcing存事件序列状态是事件回放的结果。Event Store就是存事件的地方。事件不可变只能追加不能修改删除。这是Event Sourcing的基石。packageesimport(contextencoding/jsonfmtsynctime)// Event 领域事件状态变更的原子记录typeEventstruct{IDstring// 事件唯一ID用于幂等去重StreamIDstring// 事件流ID对应聚合根ID如账户IDEventTypestring// 事件类型如MoneyDepositedData json.RawMessage// 事件数据序列化存储Versionint// 事件在流中的版本号从1递增Timestamp time.Time// 发生时间}// EventStore 事件存储接口// 所有状态变更通过Append事件记录通过Load回放typeEventStoreinterface{// Append 追加事件到流带乐观锁校验期望版本Append(ctx context.Context,streamIDstring,expectedVersionint,events[]Event)error// Load 加载流的全部事件按版本号升序Load(ctx context.Context,streamIDstring)([]Event,error)// LoadFrom 从指定版本开始加载用于分段回放LoadFrom(ctx context.Context,streamIDstring,fromVersionint)([]Event,error)}// MemoryEventStore 内存事件存储实现// 生产环境用PostgreSQL或专用EventStoreDBtypeMemoryEventStorestruct{mu sync.Mutex streamsmap[string][]Event// 流ID到事件列表}// NewMemoryEventStore 创建内存事件存储funcNewMemoryEventStore()*MemoryEventStore{returnMemoryEventStore{streams:make(map[string][]Event),}}// Append 追加事件带乐观锁// expectedVersion用于并发控制防止两个请求同时改同一聚合func(s*MemoryEventStore)Append(ctx context.Context,streamIDstring,expectedVersionint,events[]Event,)error{s.mu.Lock()defers.mu.Unlock()// 检查当前流的版本是否与期望一致current:len(s.streams[streamID])ifcurrent!expectedVersion{returnfmt.Errorf(concurrency conflict: expected version %d, got %d,expectedVersion,current)}// 追加事件版本号递增fori:rangeevents{currentevents[i].Versioncurrent events[i].StreamIDstreamID s.streams[streamID]append(s.streams[streamID],events[i])}returnnil}// Load 加载流的全部事件func(s*MemoryEventStore)Load(ctx context.Context,streamIDstring)([]Event,error){s.mu.Lock()defers.mu.Unlock()events:s.streams[streamID]// 返回副本避免外部修改result:make([]Event,len(events))copy(result,events)returnresult,nil}// LoadFrom 从指定版本开始加载func(s*MemoryEventStore)LoadFrom(ctx context.Context,streamIDstring,fromVersionint,)([]Event,error){s.mu.Lock()defers.mu.Unlock()events:s.streams[streamID]varresult[]Eventfor_,e:rangeevents{ife.VersionfromVersion{resultappend(result,e)}}returnresult,nil}乐观锁是关键。两个请求同时改同一个账户第一个成功第二个发现版本号变了直接报错让上层重试。这避免了并发覆盖。二、聚合根快照与事件回放事件多了每次都从头回放很慢。一个账户操作了一年事件可能有几万条每次算余额都要回放几万条。解决办法是快照。每隔N个事件存一份当前状态回放时从最近的快照开始只回放快照之后的事件。packageesimport(contextencoding/jsonfmtsync)// Snapshot 聚合根快照typeSnapshotstruct{StreamIDstring// 流IDVersionint// 快照对应的事件版本号State json.RawMessage// 聚合状态序列化Timestampstring// 快照时间}// SnapshotStore 快照存储接口typeSnapshotStoreinterface{// Save 保存快照Save(ctx context.Context,snap Snapshot)error// Load 加载最近的快照Load(ctx context.Context,streamIDstring)(*Snapshot,error)}// MemorySnapshotStore 内存快照存储typeMemorySnapshotStorestruct{mu sync.Mutex snapshotsmap[string]Snapshot// 流ID到快照}// NewMemorySnapshotStore 创建内存快照存储funcNewMemorySnapshotStore()*MemorySnapshotStore{returnMemorySnapshotStore{snapshots:make(map[string]Snapshot),}}// Save 保存快照覆盖旧的func(s*MemorySnapshotStore)Save(ctx context.Context,snap Snapshot)error{s.mu.Lock()defers.mu.Unlock()s.snapshots[snap.StreamID]snapreturnnil}// Load 加载最近快照func(s*MemorySnapshotStore)Load(ctx context.Context,streamIDstring)(*Snapshot,error){s.mu.Lock()defers.mu.Unlock()snap,ok:s.snapshots[streamID]if!ok{returnnil,nil// 没有快照返回nil}returnsnap,nil}// Account 账户聚合根// 状态由事件回放得到聚合根只暴露命令方法typeAccountstruct{IDstring// 账户IDBalancefloat64// 余额由事件回放计算Versionint// 当前版本号}// Apply 应用事件更新聚合状态// 这是事件回放的核心每个事件类型对应一个状态变更func(a*Account)Apply(event Event)error{switchevent.EventType{caseAccountCreated:// 反序列化事件数据vardatastruct{InitialBalancefloat64json:initial_balance}iferr:json.Unmarshal(event.Data,data);err!nil{returnerr}a.Balancedata.InitialBalancecaseMoneyDeposited:vardatastruct{Amountfloat64json:amount}iferr:json.Unmarshal(event.Data,data);err!nil{returnerr}a.Balancedata.AmountcaseMoneyWithdrawn:vardatastruct{Amountfloat64json:amount}iferr:json.Unmarshal(event.Data,data);err!nil{returnerr}a.Balance-data.Amount}// 更新版本号a.Versionevent.Versionreturnnil}// FromEvents 从事件流重建聚合状态// 有快照就从快照开始没有就从头回放func(a*Account)FromEvents(ctx context.Context,events[]Event,snapshot*Snapshot,)error{// 有快照先恢复快照状态ifsnapshot!nil{varstatestruct{Balancefloat64json:balance}iferr:json.Unmarshal(snapshot.State,state);err!nil{returnerr}a.Balancestate.Balance a.Versionsnapshot.Version}// 回放快照之后的事件for_,event:rangeevents{// 跳过快照之前的事件ifsnapshot!nilevent.Versionsnapshot.Version{continue}iferr:a.Apply(event);err!nil{returnfmt.Errorf(replay event %s failed: %w,event.EventType,err)}}returnnil}// SnapshotInterval 快照间隔每多少个事件存一次快照constSnapshotInterval100// Repository 聚合根仓储// 负责加载和保存聚合内部处理快照typeRepositorystruct{eventStore EventStore snapshotStore SnapshotStore}// Load 加载聚合自动用快照加速func(r*Repository)Load(ctx context.Context,idstring)(*Account,error){account:Account{ID:id}// 先加载快照snap,err:r.snapshotStore.Load(ctx,id)iferr!nil{returnnil,err}// 加载事件流events,err:r.eventStore.Load(ctx,id)iferr!nil{returnnil,err}// 从事件回放状态iferr:account.FromEvents(ctx,events,snap);err!nil{returnnil,err}returnaccount,nil}// Save 保存聚合追加事件并按间隔存快照func(r*Repository)Save(ctx context.Context,account*Account,newEvents[]Event)error{// 追加事件带乐观锁iferr:r.eventStore.Append(ctx,account.ID,account.Version,newEvents);err!nil{returnerr}// 每SnapshotInterval个事件存一次快照newVersion:account.Versionlen(newEvents)ifnewVersion%SnapshotInterval0{state,_:json.Marshal(struct{Balancefloat64json:balance}{Balance:account.Balance})snap:Snapshot{StreamID:account.ID,Version:newVersion,State:state,Timestamp:now,}returnr.snapshotStore.Save(ctx,snap)}returnnil}幂等处理靠事件ID去重。事件ID重复说明是重放Apply时检查ID是否已应用过。生产环境会把已应用的事件ID记下来回放时跳过。三、独家踩坑:事件版本迁移导致老数据无法回放这个坑我踩得最惨。系统上线半年后产品要求存款事件加一个渠道字段。我们直接改了事件结构新增了Channel字段。上线后老事件回放直接报错因为老事件数据里没有Channel字段反序列化失败。// 错误做法: 直接改事件结构// 老事件没有Channel字段Unmarshal报错或得到零值typeMoneyDepositedV2struct{Amountfloat64json:amountChannelstringjson:channel// 新增字段老数据没有}事件一旦写入就是不可变的结构改了老事件回放不了。正确做法是事件版本加版本号写一个upcaster把老版本事件转换成新版本。packageesimport(encoding/jsonfmt)// EventUpcaster 事件版本迁移器// 把老版本事件数据升级成新版本结构typeEventUpcasterinterface{Upcast(event Event)(Event,error)}// MoneyDepositedUpcaster 存款事件升级器// 把V1(无Channel字段)升级到V2(有Channel字段)typeMoneyDepositedUpcasterstruct{}// Upcast 把老事件升级到新版本func(u*MoneyDepositedUpcaster)Upcast(event Event)(Event,error){ifevent.Version2{// 已经是新版本直接返回returnevent,nil}// 解析老版本数据varoldstruct{Amountfloat64json:amount}iferr:json.Unmarshal(event.Data,old);err!nil{returnevent,fmt.Errorf(upcast failed: %w,err)}// 构造新版本数据补充默认值newData,_:json.Marshal(struct{Amountfloat64json:amountChannelstringjson:channel}{Amount:old.Amount,Channel:unknown,// 老数据没有渠道填默认值})event.DatanewData event.Version2returnevent,nil}// EventUpcasterRegistry 升级器注册表// 按事件类型注册upcaster回放时自动应用typeEventUpcasterRegistrystruct{upcastersmap[string]EventUpcaster}// NewEventUpcasterRegistry 创建注册表funcNewEventUpcasterRegistry()*EventUpcasterRegistry{returnEventUpcasterRegistry{upcasters:make(map[string]EventUpcaster),}}// Register 注册事件升级器func(r*EventUpcasterRegistry)Register(eventTypestring,u EventUpcaster){r.upcasters[eventType]u}// Upcast 应用升级找不到upcaster就原样返回func(r*EventUpcasterRegistry)Upcast(event Event)(Event,error){u,ok:r.upcasters[event.EventType]if!ok{returnevent,nil}returnu.Upcast(event)}经验是事件结构变更要提前规划。新字段加默认值用upcaster在回放时转换永远不要改老事件数据本身。上线前把所有老事件跑一遍upcast验证确认能正常回放。四、对比分析方案存储内容审计能力状态重建复杂度适用场景CRUD最终状态弱不可重建低简单业务Event Sourcing事件序列强可回放重建高金融审计快照状态加事件中快照加速回放中事件量大的ESCRUD只存最终状态简单高效但没有变更历史出问题查不到过程。Event Sourcing存全部事件审计能力强状态可重建代价是复杂度高每次读要回放。快照是Event Sourcing的优化定期存状态减少回放量适合事件积累多的场景。总结与预告Event Sourcing把状态变更存成不可变事件流状态由事件回放得到。Event Store负责追加和加载事件乐观锁防并发覆盖。聚合根快照每N个事件存一次加速回放。事件结构变更要用upcaster迁移老版本别直接改老数据。下一篇聊Saga模式看跨服务的分布式事务怎么用补偿事务保证最终一致。