幂等性设计模式:构建可靠的分布式系统

深入讲解幂等性在分布式系统中的重要性,涵盖幂等键、状态机、Token机制等核心模式,详解数据库唯一约束、请求去重、重试策略等实战技巧,提供完整的Go语言实现代码。

引言

在网络不可靠的分布式环境中,请求可能因为超时、重试、消息重复投递等原因被执行多次。幂等性设计确保同一操作执行一次或多次产生相同的结果,是构建可靠系统的基石。

什么是幂等性

幂等操作:无论执行多少次,结果都相同

✅ 幂等操作示例:
- UPDATE users SET email='new@example.com' WHERE id=123
- DELETE FROM orders WHERE id=456
- PUT /users/123 (完整替换)

❌ 非幂等操作示例:
- INSERT INTO logs (message) VALUES ('event')  -- 每次插入新记录
- UPDATE accounts SET balance = balance + 100 WHERE id=123  -- 余额累加
- POST /orders (创建订单)  -- 每次创建新订单

HTTP方法幂等性

方法幂等性说明
GET✅ 幂等读取操作,无副作用
PUT✅ 幂等完整替换资源
DELETE✅ 幂等删除资源(删除已删除的资源仍返回成功)
POST❌ 非幂等创建资源,需要额外处理
PATCH⚠️ 视情况部分更新,取决于实现

幂等键模式(Idempotency Key)

实现原理

type IdempotencyHandler struct {
    store IdempotencyStore
    ttl   time.Duration
}

type IdempotencyStore interface {
    // 保存幂等键和响应
    Save(ctx context.Context, key string, response *IdempotentResponse) error
    
    // 查询幂等键
    Get(ctx context.Context, key string) (*IdempotentResponse, error)
    
    // 删除过期键
    Cleanup(ctx context.Context, before time.Time) error
}

type IdempotentResponse struct {
    StatusCode int               `json:"status_code"`
    Headers    map[string]string `json:"headers"`
    Body       []byte            `json:"body"`
    CreatedAt  time.Time         `json:"created_at"`
}

func (h *IdempotencyHandler) Middleware(next http.Handler) http.Handler {
    return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
        // 仅处理POST请求(非幂等方法)
        if r.Method != http.MethodPost {
            next.ServeHTTP(w, r)
            return
        }
        
        // 提取幂等键
        idempotencyKey := r.Header.Get("Idempotency-Key")
        if idempotencyKey == "" {
            http.Error(w, "Missing Idempotency-Key header", http.StatusBadRequest)
            return
        }
        
        // 检查是否已处理过
        cachedResponse, err := h.store.Get(r.Context(), idempotencyKey)
        if err == nil && cachedResponse != nil {
            // 返回缓存的响应
            h.writeResponse(w, cachedResponse)
            return
        }
        
        // 首次处理:记录响应
        recorder := &responseRecorder{
            ResponseWriter: w,
            body:           &bytes.Buffer{},
        }
        
        next.ServeHTTP(recorder, r)
        
        // 缓存响应
        response := &IdempotentResponse{
            StatusCode: recorder.statusCode,
            Headers:    recorder.Header(),
            Body:       recorder.body.Bytes(),
            CreatedAt:  time.Now(),
        }
        
        h.store.Save(r.Context(), idempotencyKey, response)
    })
}

func (h *IdempotencyHandler) writeResponse(w http.ResponseWriter, resp *IdempotentResponse) {
    for key, value := range resp.Headers {
        w.Header().Set(key, value)
    }
    w.Header().Set("X-Idempotent-Replay", "true")
    w.WriteHeader(resp.StatusCode)
    w.Write(resp.Body)
}

Redis幂等存储

type RedisIdempotencyStore struct {
    client *redis.Client
    ttl    time.Duration
}

func NewRedisIdempotencyStore(client *redis.Client, ttl time.Duration) *RedisIdempotencyStore {
    return &RedisIdempotencyStore{
        client: client,
        ttl:    ttl,
    }
}

func (s *RedisIdempotencyStore) Save(ctx context.Context, key string, response *IdempotentResponse) error {
    data, err := json.Marshal(response)
    if err != nil {
        return err
    }
    
    return s.client.Set(ctx, "idempotency:"+key, data, s.ttl).Err()
}

func (s *RedisIdempotencyStore) Get(ctx context.Context, key string) (*IdempotentResponse, error) {
    data, err := s.client.Get(ctx, "idempotency:"+key).Bytes()
    if err == redis.Nil {
        return nil, nil
    }
    if err != nil {
        return nil, err
    }
    
    var response IdempotentResponse
    err = json.Unmarshal(data, &response)
    return &response, err
}

客户端使用示例

// 支付客户端
type PaymentClient struct {
    httpClient *http.Client
}

func (c *PaymentClient) ProcessPayment(ctx context.Context, paymentID string, amount Money) error {
    // 生成幂等键(基于业务唯一标识)
    idempotencyKey := fmt.Sprintf("payment:%s", paymentID)
    
    reqBody := map[string]interface{}{
        "payment_id": paymentID,
        "amount":     amount.Amount,
        "currency":   amount.Currency,
    }
    
    bodyBytes, _ := json.Marshal(reqBody)
    req, _ := http.NewRequestWithContext(ctx, "POST", "https://api.example.com/payments", bytes.NewReader(bodyBytes))
    req.Header.Set("Content-Type", "application/json")
    req.Header.Set("Idempotency-Key", idempotencyKey)
    
    // 重试逻辑
    for attempt := 0; attempt < 3; attempt++ {
        resp, err := c.httpClient.Do(req)
        if err != nil {
            // 网络错误,重试(使用相同的幂等键)
            time.Sleep(time.Duration(attempt) * time.Second)
            continue
        }
        
        if resp.StatusCode == http.StatusOK || resp.StatusCode == http.StatusCreated {
            return nil
        }
        
        if resp.StatusCode == http.StatusTooManyRequests {
            // 限流,等待后重试
            time.Sleep(5 * time.Second)
            continue
        }
        
        // 其他错误,不重试
        return fmt.Errorf("payment failed: %d", resp.StatusCode)
    }
    
    return errors.New("payment failed after retries")
}

数据库唯一约束

防止重复创建

-- 订单表:使用唯一约束防止重复创建
CREATE TABLE orders (
    id UUID PRIMARY KEY,
    order_number VARCHAR(50) UNIQUE NOT NULL,  -- 唯一约束
    user_id UUID NOT NULL,
    total_amount DECIMAL(10,2) NOT NULL,
    status VARCHAR(20) NOT NULL,
    created_at TIMESTAMP DEFAULT NOW()
);

-- 支付记录表
CREATE TABLE payments (
    id UUID PRIMARY KEY,
    order_id UUID NOT NULL,
    payment_provider VARCHAR(50) NOT NULL,
    provider_transaction_id VARCHAR(100),
    amount DECIMAL(10,2) NOT NULL,
    status VARCHAR(20) NOT NULL,
    created_at TIMESTAMP DEFAULT NOW(),
    
    -- 同一订单同一提供商的交易ID唯一
    UNIQUE(order_id, payment_provider, provider_transaction_id)
);
type OrderService struct {
    db *sql.DB
}

func (s *OrderService) CreateOrder(ctx context.Context, cmd CreateOrderCommand) (*Order, error) {
    orderID := uuid.New()
    
    _, err := s.db.ExecContext(ctx, `
        INSERT INTO orders (id, order_number, user_id, total_amount, status)
        VALUES ($1, $2, $3, $4, 'created')
    `, orderID, cmd.OrderNumber, cmd.UserID, cmd.TotalAmount)
    
    if err != nil {
        if isUniqueViolation(err) {
            // 订单号已存在,返回现有订单
            return s.GetOrderByNumber(ctx, cmd.OrderNumber)
        }
        return nil, err
    }
    
    return &Order{
        ID:          orderID,
        OrderNumber: cmd.OrderNumber,
        UserID:      cmd.UserID,
        TotalAmount: cmd.TotalAmount,
        Status:      "created",
    }, nil
}

func isUniqueViolation(err error) bool {
    // PostgreSQL: 23505
    // MySQL: 1062
    return strings.Contains(err.Error(), "23505") || 
           strings.Contains(err.Error(), "1062")
}

状态机模式

防止重复状态转换

type OrderStateMachine struct {
    validTransitions map[OrderStatus][]OrderStatus
}

type OrderStatus string

const (
    StatusCreated   OrderStatus = "created"
    StatusPaid      OrderStatus = "paid"
    StatusShipped   OrderStatus = "shipped"
    StatusDelivered OrderStatus = "delivered"
    StatusCancelled OrderStatus = "cancelled"
)

func NewOrderStateMachine() *OrderStateMachine {
    return &OrderStateMachine{
        validTransitions: map[OrderStatus][]OrderStatus{
            StatusCreated:   {StatusPaid, StatusCancelled},
            StatusPaid:      {StatusShipped, StatusCancelled},
            StatusShipped:   {StatusDelivered},
            StatusDelivered: {},
            StatusCancelled: {},
        },
    }
}

func (sm *OrderStateMachine) CanTransition(from, to OrderStatus) bool {
    validTargets := sm.validTransitions[from]
    for _, target := range validTargets {
        if target == to {
            return true
        }
    }
    return false
}

type OrderService struct {
    db          *sql.DB
    stateMachine *OrderStateMachine
}

func (s *OrderService) MarkAsPaid(ctx context.Context, orderID string) error {
    // 使用乐观锁和状态检查
    result, err := s.db.ExecContext(ctx, `
        UPDATE orders 
        SET status = 'paid', updated_at = NOW()
        WHERE id = $1 AND status = 'created'
    `, orderID)
    
    if err != nil {
        return err
    }
    
    rowsAffected, _ := result.RowsAffected()
    if rowsAffected == 0 {
        // 检查订单状态
        var currentStatus string
        err := s.db.QueryRowContext(ctx, 
            "SELECT status FROM orders WHERE id = $1", orderID).Scan(&currentStatus)
        
        if err != nil {
            return err
        }
        
        if currentStatus == "paid" {
            // 已经是支付状态,幂等返回成功
            return nil
        }
        
        return fmt.Errorf("invalid state transition from %s to paid", currentStatus)
    }
    
    return nil
}

Token机制(Prevention Token)

两步提交

type OrderService struct {
    db    *sql.DB
    redis *redis.Client
}

// 步骤1:创建订单Token
func (s *OrderService) CreateOrderToken(ctx context.Context, userID string) (string, error) {
    token := uuid.New().String()
    
    // 保存Token信息(5分钟过期)
    tokenData := map[string]interface{}{
        "user_id":   userID,
        "created_at": time.Now().Unix(),
        "used":      false,
    }
    
    data, _ := json.Marshal(tokenData)
    err := s.redis.Set(ctx, "order_token:"+token, data, 5*time.Minute).Err()
    if err != nil {
        return "", err
    }
    
    return token, nil
}

// 步骤2:使用Token创建订单
func (s *OrderService) CreateOrder(ctx context.Context, token string, cmd CreateOrderCommand) (*Order, error) {
    // 验证Token
    tokenData, err := s.redis.Get(ctx, "order_token:"+token).Bytes()
    if err == redis.Nil {
        return nil, errors.New("invalid or expired token")
    }
    if err != nil {
        return nil, err
    }
    
    var data map[string]interface{}
    json.Unmarshal(tokenData, &data)
    
    if data["used"].(bool) {
        return nil, errors.New("token already used")
    }
    
    // 标记Token已使用(原子操作)
    success, err := s.redis.Eval(ctx, `
        local data = redis.call('GET', KEYS[1])
        local parsed = cjson.decode(data)
        if parsed.used == false then
            parsed.used = true
            redis.call('SET', KEYS[1], cjson.encode(parsed), 'EX', 300)
            return 1
        end
        return 0
    `, []string{"order_token:" + token}).Int()
    
    if err != nil {
        return nil, err
    }
    
    if success == 0 {
        return nil, errors.New("token already used")
    }
    
    // 创建订单
    return s.doCreateOrder(ctx, cmd)
}

重试策略

指数退避重试

type RetryConfig struct {
    MaxAttempts     int
    InitialInterval time.Duration
    MaxInterval     time.Duration
    Multiplier      float64
}

func RetryWithBackoff(ctx context.Context, config RetryConfig, operation func() error) error {
    var lastErr error
    interval := config.InitialInterval
    
    for attempt := 0; attempt < config.MaxAttempts; attempt++ {
        err := operation()
        if err == nil {
            return nil
        }
        
        lastErr = err
        
        // 检查是否应该重试
        if !shouldRetry(err) {
            return err
        }
        
        // 指数退避
        select {
        case <-time.After(interval):
            interval = time.Duration(float64(interval) * config.Multiplier)
            if interval > config.MaxInterval {
                interval = config.MaxInterval
            }
        case <-ctx.Done():
            return ctx.Err()
        }
    }
    
    return fmt.Errorf("max retries exceeded: %w", lastErr)
}

func shouldRetry(err error) bool {
    // 网络错误、超时、5xx错误可以重试
    if errors.Is(err, context.DeadlineExceeded) {
        return true
    }
    
    // HTTP状态码检查
    if httpErr, ok := err.(*HTTPError); ok {
        return httpErr.StatusCode >= 500 || httpErr.StatusCode == 429
    }
    
    return false
}

// 使用示例
err := RetryWithBackoff(ctx, RetryConfig{
    MaxAttempts:     5,
    InitialInterval: 1 * time.Second,
    MaxInterval:     30 * time.Second,
    Multiplier:      2.0,
}, func() error {
    return paymentClient.ProcessPayment(ctx, payment)
})

总结

幂等性设计模式选择指南:

模式适用场景优点缺点
幂等键API请求通用性强需要客户端配合
唯一约束数据库操作数据库级别保证仅限数据库
状态机状态转换逻辑清晰需要设计状态流
Token机制表单提交防止重复提交增加复杂度

关键原则:

  • 设计API时优先考虑幂等性
  • 使用唯一标识符作为幂等键
  • 合理设置幂等键的过期时间
  • 结合重试策略处理网络故障
  • 日志记录重复请求,便于排查

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「backend」更多文章

  1. WebSocket实时通信架构:从连接到百万并发的实战指南
  2. BFF架构模式:为不同前端定制专属后端服务
  3. 蓝绿部署与金丝雀发布:零停机部署策略实战