引言
在网络不可靠的分布式环境中,请求可能因为超时、重试、消息重复投递等原因被执行多次。幂等性设计确保同一操作执行一次或多次产生相同的结果,是构建可靠系统的基石。
什么是幂等性
幂等操作:无论执行多少次,结果都相同
✅ 幂等操作示例:
- 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(¤tStatus)
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时优先考虑幂等性
- 使用唯一标识符作为幂等键
- 合理设置幂等键的过期时间
- 结合重试策略处理网络故障
- 日志记录重复请求,便于排查
延伸阅读
- Stripe: Designing Idempotent APIs
- AWS Architecture Blog: Idempotency
- Martin Kleppmann: Designing Data-Intensive Applications
- RESTful Web APIs - Idempotency
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。