引言
事件溯源(Event Sourcing)和CQRS是构建复杂业务系统的强大模式。事件溯源将状态变更记录为事件序列,CQRS将读写操作分离优化。两者结合能够构建高并发、可审计、易扩展的系统。
核心概念
事件溯源(Event Sourcing)
传统CRUD模式:
┌─────────────────────────────────────────┐
│ UPDATE orders SET status='shipped' │
│ WHERE id=123; │
│ │
│ 问题:状态被覆盖,无法追溯历史 │
└─────────────────────────────────────────┘
事件溯源模式:
┌─────────────────────────────────────────┐
│ 事件1: OrderCreated(id=123, user=456) │
│ 事件2: OrderPaid(id=123, amount=100) │
│ 事件3: OrderShipped(id=123, tracking=XYZ)│
│ │
│ 优势:完整历史、可回放、可审计 │
└─────────────────────────────────────────┘
CQRS(Command Query Responsibility Segregation)
传统模式(读写同一模型):
┌─────────┐ ┌──────────┐
│ Client │ ──────▶ │ Database │
└─────────┘ └──────────┘
CQRS模式(读写分离):
┌─────────┐ 写入 ┌──────────┐ 事件 ┌──────────┐
│ Client │ ──────▶ │ 写模型 │ ──────▶ │ 事件存储 │
└─────────┘ └──────────┘ └──────────┘
▲ │
│ 读取 │ 投影
│ ▼
┌─────────┐ ┌──────────┐ ┌──────────┐
│ Client │ ◀────── │ 读模型 │ ◀────── │ 投影服务 │
└─────────┘ └──────────┘ └──────────┘
事件存储实现
事件定义
// 领域事件基类
type DomainEvent struct {
EventID string `json:"event_id"`
AggregateID string `json:"aggregate_id"`
EventType string `json:"event_type"`
Timestamp time.Time `json:"timestamp"`
Version int `json:"version"`
Metadata map[string]string `json:"metadata,omitempty"`
Payload interface{} `json:"payload"`
}
// 订单领域事件
type OrderCreatedEvent struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Items []OrderItem `json:"items"`
TotalAmount Money `json:"total_amount"`
CreatedAt time.Time `json:"created_at"`
}
type OrderPaidEvent struct {
OrderID string `json:"order_id"`
PaymentID string `json:"payment_id"`
Amount Money `json:"amount"`
PaidAt time.Time `json:"paid_at"`
}
type OrderShippedEvent struct {
OrderID string `json:"order_id"`
TrackingNumber string `json:"tracking_number"`
Carrier string `json:"carrier"`
ShippedAt time.Time `json:"shipped_at"`
}
type OrderCancelledEvent struct {
OrderID string `json:"order_id"`
Reason string `json:"reason"`
CancelledBy string `json:"cancelled_by"`
CancelledAt time.Time `json:"cancelled_at"`
}
事件存储接口
type EventStore interface {
// 保存事件到聚合根
Save(ctx context.Context, aggregateID string, events []DomainEvent) error
// 加载聚合根的所有事件
Load(ctx context.Context, aggregateID string) ([]DomainEvent, error)
// 从特定版本开始加载
LoadFromVersion(ctx context.Context, aggregateID string, fromVersion int) ([]DomainEvent, error)
// 订阅事件流(用于投影)
Subscribe(ctx context.Context, eventTypes []string, handler EventHandler) error
}
type EventHandler func(ctx context.Context, event DomainEvent) error
PostgreSQL事件存储实现
type PostgresEventStore struct {
db *sql.DB
}
func NewPostgresEventStore(db *sql.DB) *PostgresEventStore {
return &PostgresEventStore{db: db}
}
// 创建事件表
/*
CREATE TABLE events (
event_id UUID PRIMARY KEY,
aggregate_id VARCHAR(255) NOT NULL,
event_type VARCHAR(255) NOT NULL,
version INT NOT NULL,
timestamp TIMESTAMP NOT NULL DEFAULT NOW(),
payload JSONB NOT NULL,
metadata JSONB,
UNIQUE(aggregate_id, version)
);
CREATE INDEX idx_events_aggregate ON events(aggregate_id, version);
CREATE INDEX idx_events_type ON events(event_type, timestamp);
*/
func (s *PostgresEventStore) Save(ctx context.Context, aggregateID string, events []DomainEvent) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
stmt, err := tx.PrepareContext(ctx, `
INSERT INTO events (event_id, aggregate_id, event_type, version, timestamp, payload, metadata)
VALUES ($1, $2, $3, $4, $5, $6, $7)
`)
if err != nil {
return err
}
defer stmt.Close()
for _, event := range events {
payloadJSON, _ := json.Marshal(event.Payload)
metadataJSON, _ := json.Marshal(event.Metadata)
_, err := stmt.ExecContext(ctx,
event.EventID,
event.AggregateID,
event.EventType,
event.Version,
event.Timestamp,
payloadJSON,
metadataJSON,
)
if err != nil {
// 处理版本冲突(并发写入)
if isUniqueViolation(err) {
return ErrConcurrencyConflict
}
return err
}
}
return tx.Commit()
}
func (s *PostgresEventStore) Load(ctx context.Context, aggregateID string) ([]DomainEvent, error) {
return s.LoadFromVersion(ctx, aggregateID, 0)
}
func (s *PostgresEventStore) LoadFromVersion(ctx context.Context, aggregateID string, fromVersion int) ([]DomainEvent, error) {
rows, err := s.db.QueryContext(ctx, `
SELECT event_id, aggregate_id, event_type, version, timestamp, payload, metadata
FROM events
WHERE aggregate_id = $1 AND version > $2
ORDER BY version ASC
`, aggregateID, fromVersion)
if err != nil {
return nil, err
}
defer rows.Close()
var events []DomainEvent
for rows.Next() {
var event DomainEvent
var payloadJSON, metadataJSON []byte
err := rows.Scan(
&event.EventID,
&event.AggregateID,
&event.EventType,
&event.Version,
&event.Timestamp,
&payloadJSON,
&metadataJSON,
)
if err != nil {
return nil, err
}
// 反序列化payload
event.Payload = s.deserializePayload(event.EventType, payloadJSON)
json.Unmarshal(metadataJSON, &event.Metadata)
events = append(events, event)
}
return events, nil
}
并发冲突处理
事件溯源系统中,多个请求同时操作同一聚合根时会产生并发冲突。通过乐观并发控制(Optimistic Concurrency Control) 和版本号机制解决:
// 版本号校验的 Save 方法
func (s *PostgresEventStore) SaveWithVersionCheck(
ctx context.Context,
aggregateID string,
expectedVersion int,
events []DomainEvent,
) error {
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelSerializable})
if err != nil {
return err
}
defer tx.Rollback()
// 1. 检查当前版本
var currentVersion int
err = tx.QueryRowContext(ctx,
"SELECT COALESCE(MAX(version), 0) FROM events WHERE aggregate_id = $1",
aggregateID,
).Scan(¤tVersion)
if err != nil {
return err
}
// 2. 乐观锁:版本不匹配说明有并发写入
if currentVersion != expectedVersion {
return fmt.Errorf("concurrency conflict: expected version %d, found %d",
expectedVersion, currentVersion)
}
// 3. 分配版本号并保存
for i, event := range events {
event.Version = currentVersion + i + 1
// 保存事件...
}
return tx.Commit()
}
// 客户端级别重试
func (h *OrderCommandHandler) handleWithRetry(
ctx context.Context,
cmd CreateOrderCommand,
maxRetries int,
) error {
for i := 0; i < maxRetries; i++ {
err := h.handleCreateOrder(ctx, cmd)
if err == nil {
return nil
}
// 如果是并发冲突,等待随机时间后重试
if strings.Contains(err.Error(), "concurrency conflict") {
jitter := time.Duration(rand.Intn(100)) * time.Millisecond
time.Sleep(time.Duration(i*50)*time.Millisecond + jitter)
continue
}
return err
}
return fmt.Errorf("max retries exceeded")
}
并发冲突解决策略对比:
| 策略 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
| 乐观锁 | 读多写少 | 无锁开销、吞吐量高 | 冲突时需要重试 |
| 悲观锁 | 写冲突频繁 | 避免重试、即时拒绝 | 降低并发性能 |
| 合并策略 | 最终一致性 | 高可用 | 合并逻辑复杂 |
| 最后写入获胜 | 允许覆盖 | 简单 | 可能丢失数据 |
事件版本化与演进
随着时间推移,事件结构不可避免地会演变。需要建立事件版本化以支持向后兼容:
// 事件版本化接口
type VersionedEvent interface {
EventVersion() int // 返回事件结构版本号
}
// 订单创建事件 v1(原始版本)
type OrderCreatedEventV1 struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Total float64 `json:"total"`
CreatedAt string `json:"created_at"`
}
// 订单创建事件 v2(增加了优惠码字段)
type OrderCreatedEventV2 struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Total float64 `json:"total"`
CouponCode string `json:"coupon_code,omitempty"`
DiscountRate float64 `json:"discount_rate,omitempty"`
CreatedAt string `json:"created_at"`
}
// 升级器:将 v1 事件升级为 v2
func UpgradeOrderCreatedV1ToV2(v1 OrderCreatedEventV1) OrderCreatedEventV2 {
return OrderCreatedEventV2{
OrderID: v1.OrderID,
UserID: v1.UserID,
Total: v1.Total,
CreatedAt: v1.CreatedAt,
// CouponCode 和 DiscountRate 使用默认值
}
}
// 反序列化时根据版本号选择对应结构体
func (s *PostgresEventStore) deserializePayload(eventType string, data []byte) interface{} {
switch eventType {
case "OrderCreated":
// 先检查 payload 中是否有 version 字段
var meta struct {
EventVersion int `json:"_event_version"`
}
json.Unmarshal(data, &meta)
if meta.EventVersion >= 2 {
var v2 OrderCreatedEventV2
json.Unmarshal(data, &v2)
return v2
}
var v1 OrderCreatedEventV1
json.Unmarshal(data, &v1)
return UpgradeOrderCreatedV1ToV2(v1)
}
return nil
}
聚合根实现
// 订单聚合根
type OrderAggregate struct {
ID string
UserID string
Status OrderStatus
Items []OrderItem
TotalAmount Money
Version int
uncommittedEvents []DomainEvent
}
func (o *OrderAggregate) Apply(event DomainEvent) {
switch e := event.Payload.(type) {
case OrderCreatedEvent:
o.ID = e.OrderID
o.UserID = e.UserID
o.Items = e.Items
o.TotalAmount = e.TotalAmount
o.Status = OrderStatusCreated
case OrderPaidEvent:
o.Status = OrderStatusPaid
case OrderShippedEvent:
o.Status = OrderStatusShipped
case OrderCancelledEvent:
o.Status = OrderStatusCancelled
}
o.Version = event.Version
}
// 创建订单(命令)
func (o *OrderAggregate) CreateOrder(cmd CreateOrderCommand) error {
if o.ID != "" {
return errors.New("order already exists")
}
if len(cmd.Items) == 0 {
return errors.New("order must have at least one item")
}
event := OrderCreatedEvent{
OrderID: cmd.OrderID,
UserID: cmd.UserID,
Items: cmd.Items,
TotalAmount: calculateTotal(cmd.Items),
CreatedAt: time.Now(),
}
o.recordEvent("OrderCreated", event)
return nil
}
// 支付订单
func (o *OrderAggregate) Pay(cmd PayOrderCommand) error {
if o.Status != OrderStatusCreated {
return errors.New("order cannot be paid in current status")
}
if cmd.Amount != o.TotalAmount {
return errors.New("payment amount mismatch")
}
event := OrderPaidEvent{
OrderID: o.ID,
PaymentID: cmd.PaymentID,
Amount: cmd.Amount,
PaidAt: time.Now(),
}
o.recordEvent("OrderPaid", event)
return nil
}
// 发货
func (o *OrderAggregate) Ship(cmd ShipOrderCommand) error {
if o.Status != OrderStatusPaid {
return errors.New("order must be paid before shipping")
}
event := OrderShippedEvent{
OrderID: o.ID,
TrackingNumber: cmd.TrackingNumber,
Carrier: cmd.Carrier,
ShippedAt: time.Now(),
}
o.recordEvent("OrderShipped", event)
return nil
}
func (o *OrderAggregate) recordEvent(eventType string, payload interface{}) {
o.Version++
event := DomainEvent{
EventID: uuid.New().String(),
AggregateID: o.ID,
EventType: eventType,
Version: o.Version,
Timestamp: time.Now(),
Payload: payload,
}
o.uncommittedEvents = append(o.uncommittedEvents, event)
o.Apply(event)
}
func (o *OrderAggregate) GetUncommittedEvents() []DomainEvent {
return o.uncommittedEvents
}
func (o *OrderAggregate) ClearUncommittedEvents() {
o.uncommittedEvents = nil
}
命令处理(Command Handler)
type OrderCommandHandler struct {
eventStore EventStore
eventBus EventBus
}
func (h *OrderCommandHandler) HandleCreateOrder(ctx context.Context, cmd CreateOrderCommand) error {
// 1. 加载聚合根(从事件重建状态)
aggregate := &OrderAggregate{}
events, err := h.eventStore.Load(ctx, cmd.OrderID)
if err != nil {
return err
}
for _, event := range events {
aggregate.Apply(event)
}
// 2. 执行业务逻辑
if err := aggregate.CreateOrder(cmd); err != nil {
return err
}
// 3. 保存新事件
uncommittedEvents := aggregate.GetUncommittedEvents()
if err := h.eventStore.Save(ctx, cmd.OrderID, uncommittedEvents); err != nil {
return err
}
// 4. 发布事件(异步投影)
for _, event := range uncommittedEvents {
h.eventBus.Publish(ctx, event)
}
aggregate.ClearUncommittedEvents()
return nil
}
查询模型与投影
读模型定义
// 订单查询模型(优化读取)
type OrderReadModel struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Status string `json:"status"`
TotalAmount float64 `json:"total_amount"`
Currency string `json:"currency"`
ItemCount int `json:"item_count"`
TrackingNumber string `json:"tracking_number,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
// 用户订单统计(聚合视图)
type UserOrderStats struct {
UserID string `json:"user_id"`
TotalOrders int `json:"total_orders"`
TotalSpent float64 `json:"total_spent"`
PendingOrders int `json:"pending_orders"`
LastOrderAt time.Time `json:"last_order_at"`
}
投影服务
type OrderProjectionService struct {
db *sql.DB
}
func (s *OrderProjectionService) Project(ctx context.Context, event DomainEvent) error {
switch e := event.Payload.(type) {
case OrderCreatedEvent:
return s.projectOrderCreated(ctx, e)
case OrderPaidEvent:
return s.projectOrderPaid(ctx, e)
case OrderShippedEvent:
return s.projectOrderShipped(ctx, e)
case OrderCancelledEvent:
return s.projectOrderCancelled(ctx, e)
}
return nil
}
func (s *OrderProjectionService) projectOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
_, err := s.db.ExecContext(ctx, `
INSERT INTO order_read_models (
order_id, user_id, status, total_amount, currency,
item_count, created_at, updated_at
) VALUES ($1, $2, $3, $4, $5, $6, $7, $7)
`,
event.OrderID,
event.UserID,
"created",
event.TotalAmount.Amount,
event.TotalAmount.Currency,
len(event.Items),
event.CreatedAt,
)
return err
}
func (s *OrderProjectionService) projectOrderPaid(ctx context.Context, event OrderPaidEvent) error {
_, err := s.db.ExecContext(ctx, `
UPDATE order_read_models
SET status = 'paid', updated_at = $2
WHERE order_id = $1
`, event.OrderID, event.PaidAt)
return err
}
func (s *OrderProjectionService) projectOrderShipped(ctx context.Context, event OrderShippedEvent) error {
_, err := s.db.ExecContext(ctx, `
UPDATE order_read_models
SET status = 'shipped', tracking_number = $2, updated_at = $3
WHERE order_id = $1
`, event.OrderID, event.TrackingNumber, event.ShippedAt)
return err
}
// 用户统计投影
type UserStatsProjectionService struct {
db *sql.DB
}
func (s *UserStatsProjectionService) Project(ctx context.Context, event DomainEvent) error {
switch e := event.Payload.(type) {
case OrderCreatedEvent:
return s.projectOrderCreated(ctx, e)
case OrderPaidEvent:
return s.projectOrderPaid(ctx, e)
}
return nil
}
func (s *UserStatsProjectionService) projectOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
_, err := s.db.ExecContext(ctx, `
INSERT INTO user_order_stats (user_id, total_orders, pending_orders, last_order_at)
VALUES ($1, 1, 1, $2)
ON CONFLICT (user_id) DO UPDATE SET
total_orders = user_order_stats.total_orders + 1,
pending_orders = user_order_stats.pending_orders + 1,
last_order_at = $2
`, event.UserID, event.CreatedAt)
return err
}
事件回放
重建聚合根状态
type OrderRepository struct {
eventStore EventStore
}
func (r *OrderRepository) GetByID(ctx context.Context, orderID string) (*OrderAggregate, error) {
events, err := r.eventStore.Load(ctx, orderID)
if err != nil {
return nil, err
}
if len(events) == 0 {
return nil, ErrOrderNotFound
}
// 从事件重建状态
aggregate := &OrderAggregate{}
for _, event := range events {
aggregate.Apply(event)
}
return aggregate, nil
}
完整事件回放(数据迁移)
type EventReplayer struct {
eventStore EventStore
projections []ProjectionService
}
func (r *EventReplayer) ReplayAll(ctx context.Context) error {
// 清空读模型
for _, projection := range r.projections {
projection.Reset(ctx)
}
// 加载所有事件(按时间顺序)
events, err := r.loadAllEvents(ctx)
if err != nil {
return err
}
log.Infof("Replaying %d events", len(events))
// 重放所有事件
for i, event := range events {
for _, projection := range r.projections {
if err := projection.Project(ctx, event); err != nil {
log.Errorf("Failed to project event %d: %v", i, err)
return err
}
}
if (i+1) % 10000 == 0 {
log.Infof("Replayed %d/%d events", i+1, len(events))
}
}
log.Info("Event replay completed")
return nil
}
最终一致性处理
事件总线
type EventBus interface {
Publish(ctx context.Context, event DomainEvent) error
Subscribe(ctx context.Context, eventType string, handler EventHandler) error
}
// Kafka事件总线实现
type KafkaEventBus struct {
producer *kafka.Producer
consumer *kafka.Consumer
topic string
}
func (b *KafkaEventBus) Publish(ctx context.Context, event DomainEvent) error {
eventJSON, _ := json.Marshal(event)
msg := &kafka.Message{
Topic: b.topic,
Key: []byte(event.AggregateID),
Value: eventJSON,
Headers: []kafka.Header{
{Key: "event_type", Value: []byte(event.EventType)},
},
}
return b.producer.Produce(ctx, msg, nil)
}
func (b *KafkaEventBus) Subscribe(ctx context.Context, eventType string, handler EventHandler) error {
go func() {
for {
msg, err := b.consumer.ReadMessage(ctx)
if err != nil {
log.Errorf("Failed to read message: %v", err)
continue
}
// 检查事件类型
msgEventType := getHeaderValue(msg, "event_type")
if msgEventType != eventType {
continue
}
var event DomainEvent
json.Unmarshal(msg.Value, &event)
if err := handler(ctx, event); err != nil {
log.Errorf("Failed to handle event: %v", err)
// 发送到死信队列
}
}
}()
return nil
}
总结
事件溯源与CQRS的核心价值:
事件溯源:
- 完整历史:所有状态变更都有记录
- 可审计:满足合规要求
- 时间旅行:可以回溯到任意时间点的状态
- 灵活性:可以从事件派生任意读模型
CQRS:
- 读写分离:分别优化读写性能
- 独立扩展:读模型和写模型可以独立扩展
- 简化查询:读模型针对查询优化
适用场景:
- 需要完整审计日志
- 读写负载差异大
- 业务逻辑复杂
- 需要时间旅行能力
不适用场景:
- 简单CRUD应用
- 团队缺乏经验
- 对一致性要求极高(需要强一致)
延伸阅读
- Martin Fowler: Event Sourcing
- Martin Fowler: CQRS
- Greg Young: CQRS Documents
- Event Store Documentation
- Axon Framework
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。