在构建复杂业务系统时,传统的 CRUD 架构往往难以应对高并发写入、审计追溯、状态回滚等需求。事件溯源(Event Sourcing)与命令查询职责分离(CQRS)作为领域驱动设计(DDD)的重要实践模式,为这些问题提供了优雅的解决方案。Go 语言凭借其出色的并发模型、简洁的语法和高效的运行时,成为实现事件驱动架构的理想选择。本文将从核心概念出发,逐步深入到代码实现与实战案例,带你完整掌握事件溯源与 CQRS 在 Go 中的落地实践。
一、事件溯源与 CQRS 的核心概念与关系
事件溯源是一种将系统状态的变更捕获为一系列不可变事件的持久化模式。与传统直接存储最终状态不同,事件溯源只记录发生了什么事情:订单已创建、库存已扣减、支付已完成。系统的当前状态可以通过重放这些事件来重建。这种模式天然具备完整的审计追溯能力,因为每一个状态变更的原因都被永久保存下来。
CQRS(Command Query Responsibility Segregation)是一种将数据写入(命令端)与数据读取(查询端)分离的架构模式。命令端负责处理写操作并维护业务不变式,查询端则针对读取场景进行极致优化。两种模型可以使用不同的数据结构和存储方案,各自独立扩展。
事件溯源与 CQRS 之间存在紧密的协作关系,但它们并非绑定使用。事件溯源为 CQRS 提供了天然的事件流基础——命令端每执行一次写操作就产生一个事件,查询端通过订阅这些事件来构建和更新专门的读模型。在实际项目中,CQRS 常常作为事件溯源架构的配套读取方案,因为直接从事件流中查询当前状态的代价很高,必须借助投影机制构建优化的读模型。
二、为什么 Go 适合实现事件驱动架构
Go 语言在事件驱动架构实现方面具备多项独特优势。首先,Go 的 goroutine 和 channel 提供了轻量级并发原语,能够高效地处理大量并发事件的发布与订阅。一个 goroutine 仅占用几 KB 栈空间,单机可轻松创建数十万个 goroutine,这对于需要维护大量聚合根和投影的消费者场景至关重要。
其次,Go 的强类型系统有助于在编译期捕获事件定义中的错误。领域事件作为系统中各个模块交互的契约,类型的安全性直接决定了系统的健壮性。Go 的接口机制也为聚合根的事件应用和快照序列化提供了灵活的抽象能力。
此外,Go 的标准库提供了高性能的 JSON 编解码、上下文传播(context.Context)以及 sync 包中的并发工具,这些都是实现事件总线、分布式 Saga 和并发投影所必需的基础设施。Go 的简洁语法也使得事件处理器、投影器和 Saga 协调器的代码保持清晰可维护。
最后,Go 的静态编译特性使得部署简单高效,启动速度快,非常适合作为微服务体系中事件驱动的服务节点。
三、领域事件设计:命名规范、版本控制、Schema 演化
领域事件是事件溯源系统的核心契约,其设计质量直接影响系统的可维护性和扩展性。
命名规范
领域事件名应采用过去式动词短语,准确描述业务上已经发生的事情。例如使用 OrderCreated 而非 CreateOrder,使用 InventoryDeducted 而非 DeductInventory。事件名应包含限界上下文前缀或直接放在上下文的包路径下,避免命名冲突。推荐的事件命名格式为:{领域对象}{动作}{V{版本号}},例如 OrderCreatedV1。
结构体定义
每个领域事件应包含以下元信息字段:事件 ID、聚合根 ID(用于关联到具体的业务实体)、事件发生时间、事件类型名、事件版本号。业务数据字段则应精确描述事件所携带的上下文信息。
package events
import (
"encoding/json"
"time"
)
// EventMeta 所有领域事件的公共元信息
type EventMeta struct {
EventID string `json:"event_id"`
AggregateID string `json:"aggregate_id"`
AggregateType string `json:"aggregate_type"`
EventType string `json:"event_type"`
Version int `json:"version"`
OccurredAt time.Time `json:"occurred_at"`
}
// OrderCreated 订单创建事件 V1
type OrderCreated struct {
EventMeta
UserID string `json:"user_id"`
Items []Item `json:"items"`
TotalAmount float64 `json:"total_amount"`
Address Address `json:"address"`
}
type Item struct {
ProductID string `json:"product_id"`
Quantity int `json:"quantity"`
UnitPrice float64 `json:"unit_price"`
}
type Address struct {
Province string `json:"province"`
City string `json:"city"`
Detail string `json:"detail"`
}
// OrderPaid 订单支付完成事件
type OrderPaid struct {
EventMeta
PaymentID string `json:"payment_id"`
PaidAmount float64 `json:"paid_amount"`
PaidAt time.Time `json:"paid_at"`
}
版本控制与 Schema 演化
在运行时间较长的系统中,事件 Schema 的演化不可避免。推荐的策略是:
- 向前兼容读取:在读事件时,缺失的字段使用零值填充,新增的字段使用
json.RawMessage延迟解析或直接在结构体中加入omitempty标签。 - 事件向上转换(Upcasting):在事件入库或读取管道中,注册版本转换器将旧版本事件转换为当前版本。
- 快照版本隔离:快照应包含事件版本信息,当事件 Schema 发生重大变化时,可切换聚合根的版本号并重建快照。
// EventUpcaster 事件向上转换器接口
type EventUpcaster interface {
// SupportedEventType 返回支持处理的事件类型名
SupportedEventType() string
// CanUpcast 判断给定版本是否能转换
CanUpcast(fromVersion int) bool
// Upcast 将旧版本 JSON 转换为当前版本 JSON
Upcast(fromVersion int, data []byte) ([]byte, error)
}
// UpcasterRegistry 转换器注册表
type UpcasterRegistry struct {
upcasters []EventUpcaster
}
func (r *UpcasterRegistry) Register(u EventUpcaster) {
r.upcasters = append(r.upcasters, u)
}
func (r *UpcasterRegistry) Upcast(eventType string, version int, data []byte) ([]byte, error) {
for _, u := range r.upcasters {
if u.SupportedEventType() == eventType && u.CanUpcast(version) {
return u.Upcast(version, data)
}
}
return data, nil
}
四、事件存储方案:关系型数据库 vs 专用 Event Store
事件存储是事件溯源系统的基石,它需要提供原子追加写入、按聚合根 ID 有序读取、乐观并发控制(OCC)等核心能力。
关系型数据库方案
对于大多数团队而言,使用现有的 PostgreSQL 或 MySQL 实现事件存储是最务实的选择。通过合理设计表结构,关系型数据库完全能够满足事件存储的核心需求:
CREATE TABLE event_store (
id BIGSERIAL PRIMARY KEY,
aggregate_id VARCHAR(255) NOT NULL,
aggregate_type VARCHAR(255) NOT NULL,
event_type VARCHAR(255) NOT NULL,
event_version INT NOT NULL,
event_data JSONB NOT NULL,
metadata JSONB,
sequence BIGSERIAL NOT NULL,
occurred_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE (aggregate_id, event_version)
);
CREATE INDEX idx_event_store_aggregate ON event_store(aggregate_id, event_version);
CREATE INDEX idx_event_store_sequence ON event_store(sequence);
关系型数据库的优势在于基础设施成熟、运维成本低、事务一致性强。通过 UNIQUE (aggregate_id, event_version) 约束天然实现了乐观并发控制。在需要复杂投影查询时,可以直接利用数据库的 JSONB/JSON 索引能力。
专用 Event Store(EventStoreDB)
EventStoreDB 是为事件溯源专门设计的开源数据库,支持原生的事件流语义、订阅模式、投影定义和集群复制。它提供了以下增强能力:
- 原生事件流:每个聚合根对应一个事件流,追加操作保证全局有序
- 内置订阅:支持持久订阅和竞争消费者模式,自动处理断线重连
- 内置投影:可用 JavaScript 编写服务器端投影,实时生成读模型
- 链接事件(Link Events):支持对事件进行索引和引用而不复制数据
在 Go 中连接 EventStoreDB 可通过官方 gRPC 客户端:
package main
import (
"context"
"fmt"
"log"
"github.com/EventStore/EventStore-Client-Go/v3/esdb"
)
func main() {
// 创建 EventStoreDB 连接
settings, err := esdb.ParseConnectionString("esdb://admin:changeit@localhost:2113?tls=false")
if err != nil {
log.Fatal(err)
}
client, err := esdb.NewClient(settings)
if err != nil {
log.Fatal(err)
}
defer client.Close()
// 读取某个聚合根的事件流
streamID := "order-12345"
readResult, err := client.ReadStream(
context.Background(),
streamID,
esdb.ReadStreamOptions{Direction: esdb.Forwards, From: esdb.Start{}},
100,
)
if err != nil {
log.Fatal(err)
}
defer readResult.Close()
for {
event, err := readResult.Recv()
if err != nil {
break
}
fmt.Printf("Event: %s, Data: %s\n", event.Event.EventType, string(event.Event.Data))
}
}
对于中小型项目和已有 PostgreSQL 基础设施的团队,推荐先从关系型方案起步;当事件吞吐量达到每秒数万条、对订阅延迟有严苛要求、或需要复杂的服务器端投影时,再考虑迁移到 EventStoreDB。
五、用 Go 实现事件存储库(Event Store Repository)
事件存储库是聚合根与底层存储之间的桥梁。它提供两个核心能力:保存聚合根产生的新事件(在乐观并发控制下),以及根据聚合根 ID 加载历史事件并重建状态。
聚合根接口定义
首先定义聚合根需要实现的接口,这是事件溯源框架的核心契约:
package domain
import (
"context"
"fmt"
)
// AggregateRoot 是事件溯源聚合根的接口
type AggregateRoot interface {
// AggregateID 返回聚合根的唯一标识
AggregateID() string
// AggregateType 返回聚合根类型名
AggregateType() string
// Version 返回当前版本号(未提交事件的基准版本)
Version() int
// SetVersion 设置版本号(由仓库调用)
SetVersion(int)
// ApplyEvent 将事件应用到聚合根状态(仅状态变更,无业务校验)
ApplyEvent(event DomainEvent) error
// UncommittedEvents 返回待提交的事件列表
UncommittedEvents() []DomainEvent
// ClearUncommittedEvents 清空待提交事件(保存成功后调用)
ClearUncommittedEvents()
}
// DomainEvent 领域事件接口
type DomainEvent interface {
EventMeta() EventMeta
}
// EventMeta 领域事件元信息
type EventMeta struct {
EventID string
AggregateID string
AggregateType string
EventType string
Version int
OccurredAt int64 // Unix timestamp ms
}
存储库实现
以下是一个基于 PostgreSQL 的完整事件存储库实现:
package infra
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/yourcompany/es-order/domain"
)
// EventStore 基于 PostgreSQL 的事件存储实现
type EventStore struct {
db *sql.DB
}
func NewEventStore(db *sql.DB) *EventStore {
return &EventStore{db: db}
}
// Save 保存聚合根的未提交事件,使用乐观并发控制
func (s *EventStore) Save(ctx context.Context, aggregate domain.AggregateRoot) error {
uncommitted := aggregate.UncommittedEvents()
if len(uncommitted) == 0 {
return nil
}
currentVersion := aggregate.Version()
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelSerializable})
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
// 乐观并发检查:确认聚合根当前版本与数据库一致
var latestVersion int
row := tx.QueryRowContext(ctx,
`SELECT COALESCE(MAX(event_version), 0) FROM event_store WHERE aggregate_id = $1`,
aggregate.AggregateID(),
)
if err := row.Scan(&latestVersion); err != nil {
return fmt.Errorf("check version: %w", err)
}
if latestVersion != currentVersion {
return fmt.Errorf("concurrency conflict: expected version %d, got %d",
currentVersion, latestVersion)
}
// 写入新事件
for i, evt := range uncommitted {
meta := evt.EventMeta()
meta.Version = currentVersion + i + 1
meta.OccurredAt = time.Now().UnixMilli()
data, err := json.Marshal(evt)
if err != nil {
return fmt.Errorf("marshal event: %w", err)
}
_, err = tx.ExecContext(ctx, `
INSERT INTO event_store (aggregate_id, aggregate_type, event_type, event_version, event_data, occurred_at)
VALUES ($1, $2, $3, $4, $5, to_timestamp($6::bigint / 1000.0))
`, aggregate.AggregateID(), aggregate.AggregateType(), meta.EventType, meta.Version, data, meta.OccurredAt)
if err != nil {
return fmt.Errorf("insert event: %w", err)
}
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit tx: %w", err)
}
// 更新内存中的版本号并清空未提交事件
aggregate.SetVersion(currentVersion + len(uncommitted))
aggregate.ClearUncommittedEvents()
return nil
}
// Load 加载聚合根:读取历史事件并重放
func (s *EventStore) Load(ctx context.Context, aggregate domain.AggregateRoot, aggregateID string) error {
rows, err := s.db.QueryContext(ctx, `
SELECT event_data FROM event_store
WHERE aggregate_id = $1
ORDER BY event_version ASC
`, aggregateID)
if err != nil {
return fmt.Errorf("query events: %w", err)
}
defer rows.Close()
version := 0
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return fmt.Errorf("scan event: %w", err)
}
evt, err := s.deserialize(raw)
if err != nil {
return fmt.Errorf("deserialize: %w", err)
}
if err := aggregate.ApplyEvent(evt); err != nil {
return fmt.Errorf("apply event: %w", err)
}
version++
}
aggregate.SetVersion(version)
return rows.Err()
}
// deserialize 将 JSON 反序列化为具体的领域事件(此处简化,实际可用工厂模式或反射注册)
func (s *EventStore) deserialize(data []byte) (domain.DomainEvent, error) {
// 先解析 event_type 字段,再做二次反序列化
var wrapper struct {
EventType string `json:"event_type"`
}
if err := json.Unmarshal(data, &wrapper); err != nil {
return nil, err
}
// 实际项目中使用注册表映射 event_type -> 具体类型
// 此处为示例简化
switch wrapper.EventType {
case "OrderCreated":
var evt domain.OrderCreated
if err := json.Unmarshal(data, &evt); err != nil {
return nil, err
}
return &evt, nil
// ... 其他事件类型
default:
return nil, fmt.Errorf("unknown event type: %s", wrapper.EventType)
}
}
序列化注册表改进
在生产环境中,建议使用注册表统一管理事件的序列化与反序列化:
package serialization
import (
"encoding/json"
"fmt"
"reflect"
)
// Registry 事件类型注册表
type Registry struct {
typeMap map[string]reflect.Type
}
func NewRegistry() *Registry {
return &Registry{typeMap: make(map[string]reflect.Type)}
}
func (r *Registry) Register(eventType string, prototype interface{}) {
r.typeMap[eventType] = reflect.TypeOf(prototype).Elem()
}
func (r *Registry) Deserialize(eventType string, data []byte) (interface{}, error) {
t, ok := r.typeMap[eventType]
if !ok {
return nil, fmt.Errorf("event type %s not registered", eventType)
}
ev := reflect.New(t).Interface()
if err := json.Unmarshal(data, ev); err != nil {
return nil, err
}
return ev, nil
}
六、投影(Projection)与读模型构建
事件存储中的事件流是非规范化的历史记录,直接用于查询非常低效。投影(Projection)就是将事件流转换并物化为优化查询的读模型的过程。
在 CQRS 架构中,投影器(Projector)作为事件订阅者,监听特定类型的事件并更新读模型数据库。读模型通常使用适合查询场景的数据库,如 Elasticsearch、MongoDB 或关系型数据库中的物化视图。
基本投影器实现
package projection
import (
"context"
"database/sql"
"encoding/json"
"fmt"
)
// OrderSummary 读模型:订单摘要
type OrderSummary struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Status string `json:"status"`
TotalAmount float64 `json:"total_amount"`
ItemCount int `json:"item_count"`
CreatedAt string `json:"created_at"`
LastUpdatedAt string `json:"last_updated_at"`
}
// OrderProjector 订单投影器
type OrderProjector struct {
readDB *sql.DB
}
func NewOrderProjector(readDB *sql.DB) *OrderProjector {
return &OrderProjector{readDB: readDB}
}
// HandleEvent 处理单个领域事件以更新读模型
func (p *OrderProjector) HandleEvent(ctx context.Context, eventType string, data []byte) error {
switch eventType {
case "OrderCreated":
var evt struct {
AggregateID string `json:"aggregate_id"`
UserID string `json:"user_id"`
Items []struct {
Quantity int `json:"quantity"`
} `json:"items"`
TotalAmount float64 `json:"total_amount"`
OccurredAt int64 `json:"occurred_at"`
}
if err := json.Unmarshal(data, &evt); err != nil {
return err
}
itemCount := 0
for _, item := range evt.Items {
itemCount += item.Quantity
}
_, err := p.readDB.ExecContext(ctx, `
INSERT INTO order_summary (order_id, user_id, status, total_amount, item_count, created_at, last_updated_at)
VALUES ($1, $2, 'created', $3, $4, to_timestamp($5::bigint/1000.0), to_timestamp($5::bigint/1000.0))
`, evt.AggregateID, evt.UserID, evt.TotalAmount, itemCount, evt.OccurredAt)
return err
case "OrderPaid":
var evt struct {
AggregateID string `json:"aggregate_id"`
OccurredAt int64 `json:"occurred_at"`
}
if err := json.Unmarshal(data, &evt); err != nil {
return err
}
_, err := p.readDB.ExecContext(ctx, `
UPDATE order_summary
SET status = 'paid', last_updated_at = to_timestamp($2::bigint/1000.0)
WHERE order_id = $1
`, evt.AggregateID, evt.OccurredAt)
return err
case "OrderShipped":
var evt struct {
AggregateID string `json:"aggregate_id"`
OccurredAt int64 `json:"occurred_at"`
}
if err := json.Unmarshal(data, &evt); err != nil {
return err
}
_, err := p.readDB.ExecContext(ctx, `
UPDATE order_summary
SET status = 'shipped', last_updated_at = to_timestamp($2::bigint/1000.0)
WHERE order_id = $1
`, evt.AggregateID, evt.OccurredAt)
return err
default:
return nil
}
}
投影器重放机制
系统升级或读模型结构变化时,可能需要重头重放事件来重建投影。可以通过游标批量读取事件:
// ReplayAll 从事件存储中从头重放所有事件以重建读模型
func (p *OrderProjector) ReplayAll(ctx context.Context, eventStore *infra.EventStore) error {
// 先清空读模型
if _, err := p.readDB.ExecContext(ctx, `TRUNCATE TABLE order_summary`); err != nil {
return fmt.Errorf("truncate: %w", err)
}
// 分页读取所有事件并重放
batchSize := 1000
lastSequence := 0
for {
rows, err := eventStore.QueryBySequence(ctx, lastSequence, batchSize)
if err != nil {
return fmt.Errorf("query events: %w", err)
}
count := 0
for rows.Next() {
var eventType string
var data []byte
var sequence int
if err := rows.Scan(&sequence, &eventType, &data); err != nil {
rows.Close()
return err
}
if err := p.HandleEvent(ctx, eventType, data); err != nil {
rows.Close()
return fmt.Errorf("handle event seq %d: %w", sequence, err)
}
lastSequence = sequence
count++
}
rows.Close()
if count < batchSize {
break
}
}
return nil
}
七、CQRS 命令端与查询端分离实现
在完整的 CQRS 实现中,命令端和查询端可以部署为独立的服务,使用各自的数据库,通过事件总线解耦。
命令端架构
命令端负责接收写请求、执行业务逻辑、产生事件并持久化。通常包含:命令对象、命令处理器、领域服务、事件存储库。
package app
import (
"context"
"fmt"
"time"
"github.com/yourcompany/es-order/domain"
"github.com/yourcompany/es-order/infra"
)
// CreateOrderCommand 创建订单命令
type CreateOrderCommand struct {
OrderID string
UserID string
Items []domain.Item
TotalAmount float64
Address domain.Address
}
// CommandHandler 命令处理器
type CommandHandler struct {
eventStore *infra.EventStore
eventBus infra.EventBus
}
func NewCommandHandler(es *infra.EventStore, bus infra.EventBus) *CommandHandler {
return &CommandHandler{eventStore: es, eventBus: bus}
}
// HandleCreateOrder 处理创建订单命令
func (h *CommandHandler) HandleCreateOrder(ctx context.Context, cmd CreateOrderCommand) error {
order := domain.NewOrder()
err := order.Create(cmd.OrderID, cmd.UserID, cmd.Items, cmd.TotalAmount, cmd.Address)
if err != nil {
return fmt.Errorf("domain error: %w", err)
}
// 保存事件到事件存储
if err := h.eventStore.Save(ctx, order); err != nil {
return fmt.Errorf("save events: %w", err)
}
// 发布事件到总线
for _, evt := range order.RecentEvents() {
if err := h.eventBus.Publish(ctx, evt); err != nil {
// 记录日志,发送失败不影响已持久化的一致性
// 生产环境可使用发件箱模式保证最终一致性
return fmt.Errorf("publish event: %w", err)
}
}
return nil
}
查询端架构
查询端只负责从读模型中查询数据,不涉及任何写操作或业务逻辑。
package query
import (
"context"
"database/sql"
"fmt"
)
// OrderQueryService 订单查询服务
type OrderQueryService struct {
readDB *sql.DB
}
func NewOrderQueryService(readDB *sql.DB) *OrderQueryService {
return &OrderQueryService{readDB: readDB}
}
// OrderDTO 订单数据传输对象
type OrderDTO struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Status string `json:"status"`
TotalAmount float64 `json:"total_amount"`
ItemCount int `json:"item_count"`
CreatedAt string `json:"created_at"`
}
// GetOrderByID 根据订单 ID 查询
func (s *OrderQueryService) GetOrderByID(ctx context.Context, orderID string) (*OrderDTO, error) {
row := s.readDB.QueryRowContext(ctx, `
SELECT order_id, user_id, status, total_amount, item_count, created_at
FROM order_summary
WHERE order_id = $1
`, orderID)
var dto OrderDTO
if err := row.Scan(&dto.OrderID, &dto.UserID, &dto.Status, &dto.TotalAmount, &dto.ItemCount, &dto.CreatedAt); err != nil {
if err == sql.ErrNoRows {
return nil, nil
}
return nil, err
}
return &dto, nil
}
// ListOrdersByUser 查询用户的订单列表,支持分页
func (s *OrderQueryService) ListOrdersByUser(ctx context.Context, userID string, page, pageSize int) ([]OrderDTO, int, error) {
offset := (page - 1) * pageSize
// 查询总数
var total int
if err := s.readDB.QueryRowContext(ctx,
`SELECT COUNT(*) FROM order_summary WHERE user_id = $1`, userID,
).Scan(&total); err != nil {
return nil, 0, err
}
// 查询分页数据
rows, err := s.readDB.QueryContext(ctx, `
SELECT order_id, user_id, status, total_amount, item_count, created_at
FROM order_summary
WHERE user_id = $1
ORDER BY created_at DESC
LIMIT $2 OFFSET $3
`, userID, pageSize, offset)
if err != nil {
return nil, 0, err
}
defer rows.Close()
var orders []OrderDTO
for rows.Next() {
var dto OrderDTO
if err := rows.Scan(&dto.OrderID, &dto.UserID, &dto.Status, &dto.TotalAmount, &dto.ItemCount, &dto.CreatedAt); err != nil {
return nil, 0, err
}
orders = append(orders, dto)
}
return orders, total, rows.Err()
}
HTTP 接口分离示例
package api
import (
"encoding/json"
"net/http"
"strings"
"github.com/yourcompany/es-order/app"
"github.com/yourcompany/es-order/query"
)
// CommandServer 命令端 HTTP 服务
type CommandServer struct {
cmdHandler *app.CommandHandler
}
func (s *CommandServer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPost && strings.HasPrefix(r.URL.Path, "/orders") {
var req struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
// 实际调用命令处理器...
w.WriteHeader(http.StatusAccepted)
return
}
http.NotFound(w, r)
}
// QueryServer 查询端 HTTP 服务
type QueryServer struct {
queryService *query.OrderQueryService
}
func (s *QueryServer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodGet && strings.HasPrefix(r.URL.Path, "/orders/") {
orderID := strings.TrimPrefix(r.URL.Path, "/orders/")
dto, err := s.queryService.GetOrderByID(r.Context(), orderID)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
if dto == nil {
http.NotFound(w, r)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(dto)
return
}
http.NotFound(w, r)
}
八、Saga 模式:分布式事务的协调方案
在微服务架构中,一个业务操作往往涉及多个聚合根甚至多个服务的数据变更。传统分布式事务(如两阶段提交)在性能和可用性方面存在局限。Saga 模式通过将长事务拆分为一系列本地事务,并使用补偿操作来回滚已完成的步骤,最终实现最终一致性。
在 Go 中实现 Saga 协调器:
package saga
import (
"context"
"fmt"
"sync"
"time"
)
// SagaStep Saga 中的单步定义
type SagaStep struct {
Name string
Action func(ctx context.Context) error
Compensate func(ctx context.Context) error
}
// Saga 编排器
type Saga struct {
name string
steps []SagaStep
}
func NewSaga(name string) *Saga {
return &Saga{name: name}
}
func (s *Saga) AddStep(step SagaStep) {
s.steps = append(s.steps, step)
}
// Execute 按顺序执行 Saga,任一失败时触发补偿
func (s *Saga) Execute(ctx context.Context) error {
completed := []int{}
for i, step := range s.steps {
if err := step.Action(ctx); err != nil {
// 触发已执行步骤的补偿
for j := len(completed) - 1; j >= 0; j-- {
compStep := s.steps[completed[j]]
if compStep.Compensate != nil {
// 补偿失败通常需要人工介入,记录日志即可
if compErr := compStep.Compensate(ctx); compErr != nil {
fmt.Printf("compensation failed for step %s: %v\n", compStep.Name, compErr)
}
}
}
return fmt.Errorf("saga failed at step %s: %w", step.Name, err)
}
completed = append(completed, i)
}
return nil
}
// OrderCreationSaga 订单创建 Saga:创建订单 -> 扣库存 -> 创建支付单
func BuildOrderCreationSaga(
orderSvc OrderService,
inventorySvc InventoryService,
paymentSvc PaymentService,
orderID string,
items []OrderItem,
) *Saga {
saga := NewSaga("CreateOrder")
var deductedInventory map[string]int
saga.AddStep(SagaStep{
Name: "CreateOrder",
Action: func(ctx context.Context) error {
return orderSvc.CreateDraftOrder(ctx, orderID, items)
},
Compensate: func(ctx context.Context) error {
return orderSvc.CancelOrder(ctx, orderID)
},
})
saga.AddStep(SagaStep{
Name: "DeductInventory",
Action: func(ctx context.Context) error {
result, err := inventorySvc.Deduct(ctx, items)
if err == nil {
deductedInventory = result
}
return err
},
Compensate: func(ctx context.Context) error {
if deductedInventory != nil {
return inventorySvc.Restore(ctx, deductedInventory)
}
return nil
},
})
saga.AddStep(SagaStep{
Name: "CreatePayment",
Action: func(ctx context.Context) error {
return paymentSvc.CreatePaymentOrder(ctx, orderID, calculateTotal(items))
},
Compensate: func(ctx context.Context) error {
return paymentSvc.CancelPaymentOrder(ctx, orderID)
},
})
return saga
}
type OrderService interface {
CreateDraftOrder(ctx context.Context, orderID string, items []OrderItem) error
CancelOrder(ctx context.Context, orderID string) error
}
type InventoryService interface {
Deduct(ctx context.Context, items []OrderItem) (map[string]int, error)
Restore(ctx context.Context, items map[string]int) error
}
type PaymentService interface {
CreatePaymentOrder(ctx context.Context, orderID string, amount float64) error
CancelPaymentOrder(ctx context.Context, orderID string) error
}
type OrderItem struct {
ProductID string
Quantity int
UnitPrice float64
}
func calculateTotal(items []OrderItem) float64 {
total := 0.0
for _, item := range items {
total += item.UnitPrice * float64(item.Quantity)
}
return total
}
Saga 模式的关键在于每个步骤都是独立的本地事务,失败时通过补偿操作来回滚。需要注意的是补偿操作本身也可能失败,因此必须设计补偿失败的监控和告警机制,必要时触发人工介入流程。
九、Event Bus 实现:内存 vs MQ(NATS/RabbitMQ/Kafka)
Event Bus 负责将事件从命令端传递到查询端投影器和其他订阅者。根据系统规模和部署模式,可以选择不同的实现方案。
内存事件总线(单进程内)
适合单体应用或命令端与查询端部署在同一进程内的场景:
package infra
import (
"context"
"fmt"
"sync"
)
// EventBus 事件总线接口
type EventBus interface {
Publish(ctx context.Context, event interface{}) error
Subscribe(eventType string, handler EventHandler)
}
// EventHandler 事件处理器函数类型
type EventHandler func(ctx context.Context, eventType string, data []byte) error
// InMemoryEventBus 内存事件总线实现
type InMemoryEventBus struct {
mu sync.RWMutex
handlers map[string][]EventHandler
ser EventSerializer
}
func NewInMemoryEventBus(ser EventSerializer) *InMemoryEventBus {
return &InMemoryEventBus{
handlers: make(map[string][]EventHandler),
ser: ser,
}
}
func (b *InMemoryEventBus) Publish(ctx context.Context, event interface{}) error {
b.mu.RLock()
defer b.mu.RUnlock()
eventType, data, err := b.ser.Serialize(event)
if err != nil {
return fmt.Errorf("serialize: %w", err)
}
// 异步调用处理器,避免阻塞发布方
for _, h := range b.handlers[eventType] {
go func(handler EventHandler) {
if err := handler(ctx, eventType, data); err != nil {
// 记录错误日志,生产环境应有重试和死信队列机制
fmt.Printf("handler error: %v\n", err)
}
}(h)
}
return nil
}
func (b *InMemoryEventBus) Subscribe(eventType string, handler EventHandler) {
b.mu.Lock()
defer b.mu.Unlock()
b.handlers[eventType] = append(b.handlers[eventType], handler)
}
// EventSerializer 事件序列化接口
type EventSerializer interface {
Serialize(event interface{}) (eventType string, data []byte, err error)
}
基于 NATS 的分布式事件总线
当命令端与查询端需要跨进程通信时,消息队列成为更可靠的选择。NATS JetStream 提供了轻量级、高性能的持久化消息流能力,非常适合事件溯源场景:
package infra
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/nats-io/nats.go"
)
// NATSEventBus 基于 NATS JetStream 的事件总线
type NATSEventBus struct {
js nats.JetStreamContext
streamName string
ser EventSerializer
}
func NewNATSEventBus(nc *nats.Conn, streamName string, ser EventSerializer) (*NATSEventBus, error) {
js, err := nc.JetStream()
if err != nil {
return nil, err
}
// 确保 stream 存在
_, err = js.StreamInfo(streamName)
if err == nats.ErrStreamNotFound {
_, err = js.AddStream(&nats.StreamConfig{
Name: streamName,
Subjects: []string{fmt.Sprintf("%s.*", streamName)},
Retention: nats.InterestPolicy,
})
if err != nil {
return nil, err
}
} else if err != nil {
return nil, err
}
return &NATSEventBus{js: js, streamName: streamName, ser: ser}, nil
}
func (b *NATSEventBus) Publish(ctx context.Context, event interface{}) error {
eventType, data, err := b.ser.Serialize(event)
if err != nil {
return err
}
subject := fmt.Sprintf("%s.%s", b.streamName, eventType)
_, err = b.js.Publish(subject, data, nats.Context(ctx))
return err
}
func (b *NATSEventBus) Subscribe(eventType string, handler EventHandler) error {
subject := fmt.Sprintf("%s.%s", b.streamName, eventType)
consumerName := fmt.Sprintf("consumer-%s-%d", eventType, time.Now().Unix())
_, err := b.js.Subscribe(subject, func(msg *nats.Msg) {
if err := handler(context.Background(), eventType, msg.Data); err != nil {
msg.Nak() // 否定确认,触发重试
return
}
msg.Ack()
}, nats.Durable(consumerName), nats.ManualAck())
return err
}
对于更高吞吐量和更复杂的消费者组管理需求,Kafka 是更成熟的选择;RabbitMQ 则在需要灵活路由和多种队列模式时表现优异。选择消息中间件时,应重点评估其顺序保证能力、至少一次投递保证、以及消费者组间的竞争消费机制。
十、事件回溯与状态重建性能优化
当聚合根的事件数量持续增长时,每次加载都重放所有事件将导致性能线性下降。以下是几项关键优化策略。
快照机制
定期或在事件数量达到阈值时,将聚合根的完整状态持久化为快照。加载时先读取最新快照,再重放快照之后的增量事件。
package infra
import (
"context"
"database/sql"
"encoding/json"
"fmt"
)
// SnapshotStore 快照存储
type SnapshotStore struct {
db *sql.DB
}
func NewSnapshotStore(db *sql.DB) *SnapshotStore {
return &SnapshotStore{db: db}
}
func (s *SnapshotStore) Save(ctx context.Context, aggregateID string, version int, state interface{}) error {
data, err := json.Marshal(state)
if err != nil {
return err
}
_, err = s.db.ExecContext(ctx, `
INSERT INTO snapshots (aggregate_id, version, state_data, created_at)
VALUES ($1, $2, $3, NOW())
ON CONFLICT (aggregate_id)
DO UPDATE SET version = EXCLUDED.version, state_data = EXCLUDED.state_data, created_at = EXCLUDED.created_at
`, aggregateID, version, data)
return err
}
func (s *SnapshotStore) Load(ctx context.Context, aggregateID string, target interface{}) (int, error) {
var version int
var data []byte
row := s.db.QueryRowContext(ctx,
`SELECT version, state_data FROM snapshots WHERE aggregate_id = $1`, aggregateID)
if err := row.Scan(&version, &data); err != nil {
if err == sql.ErrNoRows {
return 0, nil
}
return 0, err
}
if err := json.Unmarshal(data, target); err != nil {
return 0, err
}
return version, nil
}
在聚合根的加载流程中集成快照:
// LoadWithSnapshot 使用快照加速的聚合根加载
func (s *EventStore) LoadWithSnapshot(
ctx context.Context,
aggregate domain.AggregateRoot,
aggregateID string,
snapshotStore *SnapshotStore,
) error {
// 先加载快照
snapshotVersion, err := snapshotStore.Load(ctx, aggregateID, aggregate)
if err != nil {
return fmt.Errorf("load snapshot: %w", err)
}
// 再加载快照版本之后的增量事件
rows, err := s.db.QueryContext(ctx, `
SELECT event_data FROM event_store
WHERE aggregate_id = $1 AND event_version > $2
ORDER BY event_version ASC
`, aggregateID, snapshotVersion)
if err != nil {
return err
}
defer rows.Close()
version := snapshotVersion
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return err
}
evt, err := s.deserialize(raw)
if err != nil {
return err
}
if err := aggregate.ApplyEvent(evt); err != nil {
return err
}
version++
}
aggregate.SetVersion(version)
return rows.Err()
}
策略建议
- 快照频率:每 50-100 个事件生成一次快照,或在聚合根变更时判断
- 快照异步化:快照保存可以异步执行,不影响命令处理主路径
- 快照压缩:对大型聚合根,可考虑使用 MessagePack 或 Protobuf 替代 JSON 以减少存储和传输开销
- 读取优化:为事件存储表添加按
aggregate_id + event_version的复合索引,确保事件读取的 O(log n) 复杂度
十一、与现有 DDD 文章的衔接
在 《94-domain-driven-design》 一文中,我们系统性地介绍了领域驱动设计的核心概念,包括限界上下文、实体、值对象、聚合根、领域服务和应用服务。事件溯源与 CQRS 并非替代 DDD 的方法论,而是 DDD 在特定复杂度场景下的增强实践。
具体来说,事件溯源中的聚合根仍然是 DDD 聚合根的具象化,它维护自身的不变式并通过领域事件与其他聚合根解耦通信。命令处理器对应 DDD 中的应用服务层,负责编排领域对象完成用例。CQRS 的查询端则是对 DDD 中仓储模式的自然拆分——写操作通过事件存储库持久化变更原因,读操作则通过专门的查询服务访问物化视图。
在领域建模层面,事件风暴(Event Storming)工作坊的输出成果——领域事件列表——可以直接作为事件溯源系统的实现基础。每张大橙色的便利贴(领域事件)都将成为事件流中的一帧数据、聚合根中的一个 ApplyEvent 分支、以及投影器中的一次读模型更新。
对于已经在实践中采用 DDD 分层架构的团队,向事件溯源演进的路径可以是渐进式的:首先在不关键的聚合根上引入事件存储,保留现有的关系型写模型;随后将部分读查询迁移至基于事件的投影;最终根据业务复杂度和审计需求,决定哪些聚合适合完全事件溯源化。
十二、完整实战:订单系统的事件溯源重构
以下展示一个完整订单系统从传统 CRUD 到事件溯源的完整重构过程。
聚合根实现
package domain
import (
"fmt"
"time"
)
// Order 订单聚合根
type Order struct {
id string
userID string
status OrderStatus
items []Item
totalAmount float64
address Address
version int
uncommitted []DomainEvent
}
type OrderStatus string
const (
OrderStatusPending OrderStatus = "pending"
OrderStatusPaid OrderStatus = "paid"
OrderStatusShipped OrderStatus = "shipped"
OrderStatusCancelled OrderStatus = "cancelled"
)
func NewOrder() *Order {
return &Order{}
}
func (o *Order) AggregateID() string { return o.id }
func (o *Order) AggregateType() string { return "Order" }
func (o *Order) Version() int { return o.version }
func (o *Order) SetVersion(v int) { o.version = v }
func (o *Order) UncommittedEvents() []DomainEvent { return o.uncommitted }
func (o *Order) ClearUncommittedEvents() { o.uncommitted = nil }
// RecentEvents 返回本次加载后新产生的事件(用于发布)
func (o *Order) RecentEvents() []DomainEvent {
return o.uncommitted
}
func (o *Order) raiseEvent(event DomainEvent) {
o.uncommitted = append(o.uncommitted, event)
}
// Create 业务方法:创建订单
func (o *Order) Create(id, userID string, items []Item, total float64, addr Address) error {
if id == "" || userID == "" || len(items) == 0 {
return fmt.Errorf("invalid order parameters")
}
o.raiseEvent(&OrderCreated{
EventMeta: EventMeta{
EventID: generateID(),
AggregateID: id,
AggregateType: "Order",
EventType: "OrderCreated",
OccurredAt: time.Now().UnixMilli(),
},
UserID: userID,
Items: items,
TotalAmount: total,
Address: addr,
})
return nil
}
// Pay 业务方法:支付订单
func (o *Order) Pay(paymentID string, amount float64) error {
if o.status != OrderStatusPending {
return fmt.Errorf("order cannot be paid in status %s", o.status)
}
if amount < o.totalAmount {
return fmt.Errorf("insufficient payment amount")
}
o.raiseEvent(&OrderPaid{
EventMeta: EventMeta{
EventID: generateID(),
AggregateID: o.id,
AggregateType: "Order",
EventType: "OrderPaid",
OccurredAt: time.Now().UnixMilli(),
},
PaymentID: paymentID,
PaidAmount: amount,
PaidAt: time.Now(),
})
return nil
}
// ApplyEvent 事件应用(无业务校验,仅状态变更)
func (o *Order) ApplyEvent(event DomainEvent) error {
switch evt := event.(type) {
case *OrderCreated:
o.id = evt.AggregateID
o.userID = evt.UserID
o.items = evt.Items
o.totalAmount = evt.TotalAmount
o.address = evt.Address
o.status = OrderStatusPending
case *OrderPaid:
o.status = OrderStatusPaid
default:
return fmt.Errorf("unknown event type")
}
return nil
}
func generateID() string {
// 简化的 ID 生成,实际应使用 UUID 或雪花算法
return fmt.Sprintf("evt-%d", time.Now().UnixNano())
}
应用服务层
package app
import (
"context"
"fmt"
"github.com/yourcompany/es-order/domain"
"github.com/yourcompany/es-order/infra"
)
type OrderAppService struct {
eventStore *infra.EventStore
snapshotStore *infra.SnapshotStore
eventBus infra.EventBus
}
func NewOrderAppService(es *infra.EventStore, ss *infra.SnapshotStore, bus infra.EventBus) *OrderAppService {
return &OrderAppService{eventStore: es, snapshotStore: ss, eventBus: bus}
}
// CreateOrder 创建订单用例
func (s *OrderAppService) CreateOrder(ctx context.Context, cmd CreateOrderCommand) (*domain.Order, error) {
order := domain.NewOrder()
if err := order.Create(cmd.OrderID, cmd.UserID, cmd.Items, cmd.TotalAmount, cmd.Address); err != nil {
return nil, fmt.Errorf("domain validation failed: %w", err)
}
if err := s.eventStore.Save(ctx, order); err != nil {
return nil, fmt.Errorf("persist failed: %w", err)
}
// 异步发布事件
for _, evt := range order.RecentEvents() {
_ = s.eventBus.Publish(ctx, evt)
}
return order, nil
}
// GetOrder 加载订单(带快照优化)
func (s *OrderAppService) GetOrder(ctx context.Context, orderID string) (*domain.Order, error) {
order := domain.NewOrder()
if err := s.eventStore.LoadWithSnapshot(ctx, order, orderID, s.snapshotStore); err != nil {
return nil, err
}
if order.AggregateID() == "" {
return nil, fmt.Errorf("order not found")
}
// 判断是否需要生成快照
const SnapshotThreshold = 50
if order.Version()-s.lastSnapshotVersion(orderID) >= SnapshotThreshold {
go func() {
_ = s.snapshotStore.Save(context.Background(), orderID, order.Version(), order)
}()
}
return order, nil
}
func (s *OrderAppService) lastSnapshotVersion(orderID string) int {
// 简化处理,实际应从缓存或表中查询
return 0
}
读模型投影器集成
package projection
import (
"context"
"encoding/json"
"fmt"
"github.com/yourcompany/es-order/infra"
)
type OrderProjectionService struct {
projector *OrderProjector
}
func NewOrderProjectionService(projector *OrderProjector) *OrderProjectionService {
return &OrderProjectionService{projector: projector}
}
// StartConsuming 启动事件消费
func (s *OrderProjectionService) StartConsuming(ctx context.Context, bus infra.EventBus) error {
bus.Subscribe("OrderCreated", func(ctx context.Context, eventType string, data []byte) error {
return s.projector.HandleEvent(ctx, eventType, data)
})
bus.Subscribe("OrderPaid", func(ctx context.Context, eventType string, data []byte) error {
return s.projector.HandleEvent(ctx, eventType, data)
})
bus.Subscribe("OrderShipped", func(ctx context.Context, eventType string, data []byte) error {
return s.projector.HandleEvent(ctx, eventType, data)
})
<-ctx.Done()
return ctx.Err()
}
启动入口
package main
import (
"context"
"database/sql"
"log"
"net/http"
"os"
"os/signal"
"syscall"
_ "github.com/lib/pq"
"github.com/yourcompany/es-order/api"
"github.com/yourcompany/es-order/app"
"github.com/yourcompany/es-order/infra"
"github.com/yourcompany/es-order/projection"
"github.com/yourcompany/es-order/query"
)
func main() {
writeDB, err := sql.Open("postgres", "postgres://user:pass@localhost/write_db?sslmode=disable")
if err != nil {
log.Fatal(err)
}
readDB, err := sql.Open("postgres", "postgres://user:pass@localhost/read_db?sslmode=disable")
if err != nil {
log.Fatal(err)
}
eventStore := infra.NewEventStore(writeDB)
snapshotStore := infra.NewSnapshotStore(writeDB)
eventBus := infra.NewInMemoryEventBus(&jsonSerializer{})
orderService := app.NewOrderAppService(eventStore, snapshotStore, eventBus)
orderProjector := projection.NewOrderProjector(readDB)
projectionService := projection.NewOrderProjectionService(orderProjector)
// 启动投影消费
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go func() {
if err := projectionService.StartConsuming(ctx, eventBus); err != nil {
log.Printf("projection stopped: %v", err)
}
}()
// 启动 API 服务
mux := http.NewServeMux()
cmdServer := &api.CommandServer{}
queryServer := &api.QueryServer{}
mux.Handle("/cmd/", http.StripPrefix("/cmd", cmdServer))
mux.Handle("/query/", http.StripPrefix("/query", queryServer))
srv := &http.Server{Addr: ":8080", Handler: mux}
go func() {
log.Println("Server listening on :8080")
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatal(err)
}
}()
// 优雅关闭
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
cancel()
if err := srv.Shutdown(context.Background()); err != nil {
log.Printf("shutdown error: %v", err)
}
}
type jsonSerializer struct{}
func (j *jsonSerializer) Serialize(event interface{}) (string, []byte, error) {
data, err := json.Marshal(event)
if err != nil {
return "", nil, err
}
var wrapper struct {
EventType string `json:"event_type"`
}
_ = json.Unmarshal(data, &wrapper)
return wrapper.EventType, data, nil
}
十三、总结
事件溯源与 CQRS 为复杂业务系统提供了强大的架构工具:不可变的事件流带来了完整的审计能力,命令与查询的分离使读写可以独立优化,Saga 模式在保持最终一致性的同时实现了跨聚合的复杂业务流程。
然而,引入这些模式也意味着系统复杂度的显著提升。团队需要权衡以下成本:开发人员的学习曲线、事件 Schema 的长期演化治理、投影延迟带来的最终一致性体验、以及运维多个数据存储的代价。对于大多数业务系统而言,并非所有聚合都适合事件溯源,推荐仅在需要强审计、频繁状态回滚或复杂跨聚合协作的核心领域引入。
Go 语言凭借其简洁高效的并发模型和类型安全的编译期检查,为落地事件驱动架构提供了坚实的技术基础。结合 PostgreSQL 的 JSONB 能力、成熟的消息中间件以及清晰的分层设计,团队可以构建出既具备业务表达力又拥有高可运维性的事件溯源系统。
在实践过程中,建议遵循渐进式演进策略:从单个聚合根的事件存储试点开始,逐步完善投影和读模型,最后根据实际负载选择合适的事件总线方案。始终保持对领域模型的关注——技术实现应为业务表达服务,而非让业务逻辑迁就技术框架。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。