go语言中cqrs需手动建模command与event:command是状态变更请求,仅校验+写db;event是已变更事实,不可变且不触发新command;es投递须通过outbox表+补偿重试,id用业务主键,时间戳由写模型透传。

Go 语言里直接用 CQRS 模式本身没有现成的“框架级支持”,它只是种架构分层思路;真正落地时,command 和 event 的协作必须靠你手动建模、调度和投递——否则容易变成伪 CQRS:写模型和读模型还是强耦合,事件不幂等、不保序、不重试,ES 最终查不到数据。
如何定义 command 与 event 的边界
很多人一上来就建 CreateUserCommand 和 UserCreated,但没想清楚谁该负责持久化、谁该负责通知。关键判断点就一个:command 是状态变更的请求,event 是状态已变更的事实。
-
command不该包含任何副作用逻辑(比如发邮件、写 ES),只做校验 + 写 DB(或事务性存储) -
event必须是不可变结构体,字段全为public readonly或导出字段+无 setter,避免监听器意外篡改 - 一个 command 可能触发多个 event(如
OrderPlaced→InventoryReserved+NotificationQueued),但 event 不得再触发新 command(否则链路失控) - 示例中常见错误:
UpdateUserCommand直接调用es.Index()—— 这破坏了命令端纯净性,也绕过了事件重放能力
event 投递到 Elasticsearch 的可靠性保障
Go 里用 elasticsearch.Client 写入 ES 很简单,但生产环境里最常崩在“写一半失败”:DB 提交成功了,event 发出去了,但 ES 写挂了,读模型就永久缺失。
在 Golang 中使用 samber/hot 进行内存缓存,支持 LRU、LFU、TinyLFU、W‑TinyLFU、S3FIFO、ARC、TwoQueue、SIEVE、FIFO 等淘汰算法,提供 TTL、缓存加载器及分片功能。
- 必须把 event 投递纳入事务外补偿流程:先落库(如 PostgreSQL 的
outbox表),再由独立 worker 轮询投递,失败则重试 + 记录死信 - 不要依赖 Kafka → ES 的直连管道(如 Logstash 或 Kafka Connect),它不保证 exactly-once,且无法按业务语义控制重试粒度
- ES 写入时设
Refresh: "wait_for",避免IndexRequest返回成功但文档不可查;同时配好RetryOnStatus和指数退避(如503时 sleep 100ms→200ms→400ms) - ES 文档 ID 强烈建议用业务主键(如
order_id),而非自增 ID 或 UUID —— 否则更新事件会生成新文档,聚合查询失效
读模型(ES)与写模型(PostgreSQL)的一致性陷阱
很多人以为“监听 event 更新 ES”就自动实现了最终一致性,结果上线后发现搜索结果滞后几分钟、重复、甚至漏数据。
- ES 默认
refresh_interval是 1s,但网络抖动或 bulk 队列积压会导致延迟放大;监控要盯elasticsearch_indices_refresh_total和elasticsearch_thread_pool_bulk_rejected - PostgreSQL 的
outbox表必须加唯一索引(event_type, aggregate_id, version),防重复投递;ES 端也要用version_type=external+ 业务版本号做乐观并发控制 - 别在 event handler 里做复杂计算(如 join 多张表生成 view)—— 这会让投递变慢、卡住整个 worker;应提前在写模型侧物化好所需字段,event 只带必要 payload
- 最隐蔽的坑:
time.Now()在 event 构造时调用,但投递可能延迟几秒,导致 ES 中created_at比 DB 中晚 —— 正确做法是让写模型生成时间戳并随 event 透传
真正的难点不在代码怎么写,而在于每个 event 的语义是否经得起重放、每个 command 的边界是否拦得住副作用、每个 ES 写入是否扛得住节点临时失联。这些地方一旦松动,CQRS 就退化成一堆异步函数调用。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










