免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

Go CQRS架构:命令查询职责分离

Go CQRS架构:命令查询职责分离 Go CQRS架构:命令查询职责分离摘要: 本篇讲解Go语言CQRS架构实现分离Command写模型和Query读模型写模型用事件驱动持久化事件流读模型用物化视图加速查询结合Event Sourcing实现状态回放分享最终一致性导致读取延迟的踩坑经验对比CQRS、传统CRUD、CQRS加Event Sourcing三种方案。开篇故事去年做电商系统重构订单表涨到2000万行。写订单还好一秒几百笔没问题。查询慢得要命运营后台按用户、按时间、按状态各种维度查一个订单列表查询要3秒。加了四个索引写入开始变慢查询也就快了30%。DBA建议做读写分离主库写从库读。从库延迟几百毫秒运营能接受。有个统计报表查不到刚下的订单运营投诉。加了强制读主库的逻辑主库压力又上来了。后来我把读写彻底拆开写走Command模型只管订单状态变更读走Query模型专门为查询优化存储结构。这就是CQRSCommand Query Responsibility Segregation。这篇把Go里的实现写清楚。一、Command和Query分离模型CQRS的核心是把写和读拆成两套模型。写模型负责业务状态变更追求正确性。读模型负责查询追求性能。两边各自优化互不干扰。先看Command这一侧。Command是一个意图表达要做什么不返回数据。packagecqrsimport(contextfmttime)// Command 命令接口所有命令实现这个接口typeCommandinterface{CommandType()string// 命令类型标识用于分发}// CreateOrderCommand 创建订单命令typeCreateOrderCommandstruct{OrderIDstring// 订单IDUserIDstring// 用户IDProductIDstring// 商品IDQuantityint// 数量Amountfloat64// 金额OccurTime time.Time// 发生时间}// CommandType 返回命令类型func(c CreateOrderCommand)CommandType()string{returnCreateOrder}// CancelOrderCommand 取消订单命令typeCancelOrderCommandstruct{OrderIDstring// 订单IDReasonstring// 取消原因OccurTime time.Time// 发生时间}// CommandType 返回命令类型func(c CancelOrderCommand)CommandType()string{returnCancelOrder}// CommandHandler 命令处理器接口// 每种命令对应一个handler负责执行业务逻辑typeCommandHandlerinterface{Handle(ctx context.Context,cmd Command)error}// CommandBus 命令总线分发命令到对应handlertypeCommandBusstruct{handlersmap[string]CommandHandler// 按类型注册的handler表}// NewCommandBus 创建命令总线funcNewCommandBus()*CommandBus{returnCommandBus{handlers:make(map[string]CommandHandler),}}// Register 注册命令处理器func(b*CommandBus)Register(cmdTypestring,h CommandHandler){b.handlers[cmdType]h}// Dispatch 分发命令到对应handler执行func(b*CommandBus)Dispatch(ctx context.Context,cmd Command)error{h,ok:b.handlers[cmd.CommandType()]if!ok{// 没有注册对应handler直接报错returnfmt.Errorf(no handler for command %s,cmd.CommandType())}returnh.Handle(ctx,cmd)}CommandBus负责分发命令到对应handler。业务代码只管发命令不管谁来处理。加新命令类型时老代码不用改符合开闭原则。Query这一侧简单很多。Query直接返回数据不做状态变更。packagecqrsimportcontext// OrderView 读模型的订单视图// 这个结构专门为查询优化字段都是查询需要的typeOrderViewstruct{OrderIDstring// 订单IDUserIDstring// 用户ID用于按用户查ProductIDstring// 商品ID用于按商品统计Statusstring// 订单状态Amountfloat64// 金额CreatedAtstring// 创建时间字符串方便排序展示}// QueryService 查询服务接口typeQueryServiceinterface{// GetByID 按订单ID查GetByID(ctx context.Context,orderIDstring)(*OrderView,error)// ListByUser 按用户分页查ListByUser(ctx context.Context,userIDstring,page,sizeint)([]*OrderView,error)// ListByStatus 按状态查ListByStatus(ctx context.Context,statusstring,page,sizeint)([]*OrderView,error)}读模型不用跟写模型用同一张表。读模型可以是一张扁平的宽表也可以是Redis缓存甚至可以是Elasticsearch索引。怎么快怎么来。二、事件驱动写模型与物化视图CQRS单独用价值有限配上事件驱动才有意思。写模型只产生事件状态表交给事件消费方更新。读模型订阅事件更新物化视图。packagecqrsimport(contextencoding/jsonfmtlogtime)// Event 领域事件写模型产生的变更记录typeEventstruct{IDstring// 事件ID用于幂等去重EventTypestring// 事件类型如OrderCreatedAggregatestring// 聚合根ID如订单IDData json.RawMessage// 事件数据序列化存储OccurTime time.Time// 发生时间Versionint// 事件版本号用于版本迁移}// EventBus 事件总线发布事件给订阅者typeEventBusinterface{Publish(ctx context.Context,event Event)error}// EventStore 事件存储持久化事件流// Event Sourcing的核心所有状态变更都存成事件typeEventStoreinterface{Append(ctx context.Context,events[]Event)errorLoad(ctx context.Context,aggregateIDstring)([]Event,error)}// OrderCommandHandler 订单命令处理器// 处理命令后产生事件不直接改状态表typeOrderCommandHandlerstruct{eventStore EventStore// 事件存储eventBus EventBus// 事件总线通知读模型}// Handle 处理创建订单命令// 流程: 校验 - 产生事件 - 持久化事件 - 发布事件func(h*OrderCommandHandler)Handle(ctx context.Context,cmd Command)error{c,ok:cmd.(CreateOrderCommand)if!ok{returnfmt.Errorf(invalid command type, want CreateOrder)}// 产生订单创建事件event:Event{ID:fmt.Sprintf(evt-%s,c.OrderID),EventType:OrderCreated,Aggregate:c.OrderID,OccurTime:c.OccurTime,Version:1,}// 把命令数据序列化进事件data,err:json.Marshal(c)iferr!nil{returnfmt.Errorf(marshal event data failed: %w,err)}event.Datadata// 先持久化事件保证事件不丢// 这是Event Sourcing的精髓事件即真相iferr:h.eventStore.Append(ctx,[]Event{event});err!nil{returnfmt.Errorf(persist event failed: %w,err)}// 发布事件通知读模型更新// 失败只记日志不影响写模型返回成功iferr:h.eventBus.Publish(ctx,event);err!nil{log.Printf(publish event failed: %v, event will be replayed later,err)}returnnil}读模型订阅事件更新物化视图。这里用内存map模拟实际会落到数据库或缓存。packagecqrsimport(contextencoding/jsonfmtsync)// OrderViewStore 读模型存储物化视图// 监听事件流实时更新typeOrderViewStorestruct{mu sync.RWMutex viewsmap[string]*OrderView// 按订单ID索引byUsermap[string][]string// 用户到订单ID列表的索引}// NewOrderViewStore 创建读模型存储funcNewOrderViewStore()*OrderViewStore{returnOrderViewStore{views:make(map[string]*OrderView),byUser:make(map[string][]string),}}// OnEvent 处理事件更新物化视图// 这就是读模型的事件消费逻辑func(s*OrderViewStore)OnEvent(ctx context.Context,event Event)error{switchevent.EventType{caseOrderCreated:// 反序列化事件数据varcmd CreateOrderCommandiferr:json.Unmarshal(event.Data,cmd);err!nil{returnerr}// 构造读模型视图view:OrderView{OrderID:cmd.OrderID,UserID:cmd.UserID,ProductID:cmd.ProductID,Status:created,Amount:cmd.Amount,CreatedAt:cmd.OccurTime.Format(2006-01-02 15:04:05),}// 更新物化视图双索引s.mu.Lock()s.views[cmd.OrderID]view s.byUser[cmd.UserID]append(s.byUser[cmd.UserID],cmd.OrderID)s.mu.Unlock()}returnnil}// GetByID 按ID查func(s*OrderViewStore)GetByID(ctx context.Context,orderIDstring)(*OrderView,error){s.mu.RLock()defers.mu.RUnlock()view,ok:s.views[orderID]if!ok{returnnil,fmt.Errorf(order %s not found,orderID)}returnview,nil}写模型产生事件读模型消费事件更新视图两边通过事件总线解耦。写模型可以慢慢做业务校验读模型可以用最快的方式存数据。配合Event Sourcing事件流就是真相之源状态丢了可以从事件回放重建。三、独家踩坑:最终一致性导致读取延迟这个坑踩得印象深刻。上线CQRS后压测创建订单接口返回成功前端立刻查订单详情查不到报错order not found。用户以为下单失败重复提交结果下了两单。原因在于写模型发完事件就返回了读模型消费事件更新视图有延迟。我们用NATS做事件总线正常延迟几毫秒高峰期消息堆积延迟到了500毫秒。前端创建成功即查的逻辑直接暴露了最终一致性问题。解决方案有两个。第一前端创建订单后轮询查询最多重试3次每次间隔200毫秒。第二产品层面接受延迟创建后跳转到订单列表而不是详情给读模型留同步时间。packagecqrsimport(contextfmttime)// ConsistentReader 最终一致性读包装// 创建后立即查的场景用带重试的读取typeConsistentReaderstruct{viewStore*OrderViewStore eventStore EventStore}// GetWithRetry 带重试的读取// 读不到就等一会再读最多重试3次func(r*ConsistentReader)GetWithRetry(ctx context.Context,orderIDstring,)(*OrderView,error){// 最多重试3次每次间隔200毫秒fori:0;i3;i{view,err:r.viewStore.GetByID(ctx,orderID)iferrnil{returnview,nil// 读到了}// 等待读模型追上事件流select{case-time.After(200*time.Millisecond):case-ctx.Done():returnnil,ctx.Err()}}returnnil,fmt.Errorf(order %s not ready, try later,orderID)}经验是CQRS系统要在产品层面接受最终一致性。关键接口可以同步双写读模型兜底或者前端做轮询。绝对不能假设写成功后读模型立即可见。四、对比分析方案读写分离查询性能复杂度一致性适用场景传统CRUD读写同模型中低强一致中小系统CQRS读写模型分离高中最终一致读写悬殊CQRSEvent Sourcing读写分离加事件流高高最终一致需审计回溯传统CRUD简单直接读写用一张表适合中小系统。CQRS读写分离各自优化查询性能好代价是两套模型要同步。CQRS加Event Sourcing多了事件流能回放重建状态复杂度也最高适合对审计和回溯有要求的场景。总结与预告CQRS把读写拆开写模型管正确性读模型管性能。事件驱动让两边解耦配合Event Sourcing还能回放重建状态。最终一致性是CQRS的代价产品层面要留出同步窗口。下一篇聊Event Sourcing看事件流怎么存、聚合根快照怎么做、老事件版本怎么迁移。
返回列表