Python幂等工具案例如何封装接口幂等

wen python案例 28

本文目录导读:

Python幂等工具案例如何封装接口幂等

  1. 基于装饰器的幂等实现
  2. 基于中间件的幂等实现
  3. 数据库级别的幂等实现
  4. 完整的幂等工具类
  5. 使用建议

我来详细介绍Python中接口幂等性的封装实现方案。

基于装饰器的幂等实现

import hashlib
import json
import time
from functools import wraps
from typing import Any, Callable, Dict, Optional
import redis
from flask import request, current_app
class IdempotentDecorator:
    """幂等性装饰器类"""
    def __init__(self, redis_client, expire_time=3600):
        """
        :param redis_client: Redis客户端实例
        :param expire_time: 幂等key过期时间(秒)
        """
        self.redis = redis_client
        self.expire_time = expire_time
    def generate_key(self, func_name: str, args: tuple, kwargs: dict) -> str:
        """生成幂等key"""
        # 获取请求的唯一标识(可以从请求头或参数中获取)
        idempotent_key = kwargs.get('idempotent_key') or request.headers.get('Idempotent-Key')
        if not idempotent_key:
            # 如果没有提供幂等键,使用参数生成
            content = f"{func_name}:{json.dumps(args)}:{json.dumps(kwargs, sort_keys=True)}"
            idempotent_key = hashlib.md5(content.encode()).hexdigest()
        return f"idempotent:{func_name}:{idempotent_key}"
    def __call__(self, func: Callable) -> Callable:
        @wraps(func)
        def wrapper(*args, **kwargs):
            # 生成幂等key
            key = self.generate_key(func.__name__, args, kwargs)
            # 检查是否已处理
            existing_result = self.redis.get(key)
            if existing_result:
                return json.loads(existing_result)
            # 设置锁防止并发
            lock_key = f"{key}:lock"
            if not self.redis.setnx(lock_key, 1):
                # 如果获取锁失败,说明请求正在处理
                return {"error": "Request is being processed"}, 409
            # 设置锁过期时间(防止死锁)
            self.redis.expire(lock_key, 10)
            try:
                # 执行原函数
                result = func(*args, **kwargs)
                # 存储结果
                self.redis.setex(key, self.expire_time, json.dumps(result))
                return result
            finally:
                # 释放锁
                self.redis.delete(lock_key)
        return wrapper
# 使用示例
idempotent_decorator = IdempotentDecorator(redis_client)
@app.route('/api/order/create', methods=['POST'])
@idempotent_decorator
def create_order():
    data = request.get_json()
    # 处理订单逻辑
    return {"order_id": "123", "status": "success"}

基于中间件的幂等实现

from flask import Flask, request, g, jsonify
import uuid
import hashlib
from datetime import datetime
class IdempotencyMiddleware:
    """幂等中间件"""
    def __init__(self, app=None, storage=None, expire_time=3600):
        self.app = app
        self.storage = storage or {}
        self.expire_time = expire_time
        if app is not None:
            self.init_app(app)
    def init_app(self, app):
        app.before_request(self.before_request)
        app.after_request(self.after_request)
    def get_idempotent_key(self):
        """从请求中提取幂等键"""
        # 优先从请求头获取
        key = request.headers.get('X-Idempotent-Key')
        if not key:
            # 从请求体获取
            if request.is_json:
                key = request.json.get('idempotent_key')
        if not key:
            # 生成基于请求的唯一标识
            content = f"{request.method}:{request.path}:{request.data}"
            key = hashlib.md5(content.encode()).hexdigest()
        return key
    def before_request(self):
        """请求前检查幂等性"""
        if request.method not in ['POST', 'PUT', 'PATCH']:
            return None
        idempotent_key = self.get_idempotent_key()
        # 检查是否已处理
        if idempotent_key in self.storage:
            existing_result = self.storage[idempotent_key]
            # 如果还在处理中,返回处理中状态
            if existing_result.get('status') == 'processing':
                return jsonify({"error": "Request is being processed"}), 409
            # 如果已处理完成,返回缓存结果
            return jsonify(existing_result['data']), 200
        # 标记为处理中
        self.storage[idempotent_key] = {
            'status': 'processing',
            'timestamp': datetime.now()
        }
        # 将key存入request上下文
        g.idempotent_key = idempotent_key
    def after_request(self, response):
        """请求后存储结果"""
        if hasattr(g, 'idempotent_key'):
            key = g.idempotent_key
            # 检查是否是错误响应
            if response.status_code >= 400:
                # 错误响应不缓存,删除处理标记
                self.storage.pop(key, None)
            else:
                # 缓存成功响应
                self.storage[key] = {
                    'status': 'completed',
                    'data': response.get_json(),
                    'timestamp': datetime.now()
                }
        return response
# 使用示例
app = Flask(__name__)
storage = {}  # 实际使用Redis
middleware = IdempotencyMiddleware(app, storage)
@app.route('/api/payment', methods=['POST'])
def make_payment():
    payment_data = request.get_json()
    # 处理支付逻辑
    return jsonify({"payment_id": "pay_123", "status": "success"})

数据库级别的幂等实现

from sqlalchemy import Column, String, DateTime, Text, create_engine
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
from contextlib import contextmanager
import json
from datetime import datetime
Base = declarative_base()
class IdempotentRecord(Base):
    """幂等记录表"""
    __tablename__ = 'idempotent_records'
    idempotent_key = Column(String(64), primary_key=True)
    request_data = Column(Text)
    response_data = Column(Text)
    status = Column(String(20))  # processing, completed, failed
    created_at = Column(DateTime, default=datetime.utcnow)
    updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
class DatabaseIdempotentTool:
    """基于数据库的幂等工具"""
    def __init__(self, db_session_factory):
        self.session_factory = db_session_factory
    @contextmanager
    def get_session(self):
        session = self.session_factory()
        try:
            yield session
            session.commit()
        except:
            session.rollback()
            raise
        finally:
            session.close()
    def check_and_process(self, idempotent_key: str, request_data: dict, 
                         process_func: callable) -> tuple:
        """检查并处理幂等请求"""
        with self.get_session() as session:
            # 查询是否存在记录
            record = session.query(IdempotentRecord).filter_by(
                idempotent_key=idempotent_key
            ).first()
            if record:
                if record.status == 'processing':
                    return {'error': '请求正在处理中'}, 409
                elif record.status == 'completed':
                    return json.loads(record.response_data), 200
                else:  # failed
                    # 删除失败记录,允许重试
                    session.delete(record)
                    session.flush()
            # 创建新记录
            record = IdempotentRecord(
                idempotent_key=idempotent_key,
                request_data=json.dumps(request_data),
                status='processing'
            )
            session.add(record)
        try:
            # 执行业务逻辑
            result = process_func(request_data)
            # 更新记录为完成状态
            with self.get_session() as session:
                record = session.query(IdempotentRecord).filter_by(
                    idempotent_key=idempotent_key
                ).first()
                record.status = 'completed'
                record.response_data = json.dumps(result)
            return result, 200
        except Exception as e:
            # 更新记录为失败状态
            with self.get_session() as session:
                record = session.query(IdempotentRecord).filter_by(
                    idempotent_key=idempotent_key
                ).first()
                record.status = 'failed'
            raise e
# 使用示例
engine = create_engine('sqlite:///idempotent.db')
Base.metadata.create_all(engine)
Session = sessionmaker(bind=engine)
tool = DatabaseIdempotentTool(Session)
@app.route('/api/transfer', methods=['POST'])
def transfer():
    data = request.get_json()
    idempotent_key = request.headers.get('X-Idempotent-Key')
    def process_transfer(transfer_data):
        # 资金转账逻辑
        return {"transfer_id": "transfer_123", "amount": transfer_data['amount']}
    result, status_code = tool.check_and_process(idempotent_key, data, process_transfer)
    return jsonify(result), status_code

完整的幂等工具类

import hashlib
import json
import time
import threading
from abc import ABC, abstractmethod
from typing import Any, Dict, Optional, Tuple
from datetime import datetime, timedelta
import redis
from functools import wraps
class IdempotentStorage(ABC):
    """幂等存储抽象类"""
    @abstractmethod
    def get(self, key: str) -> Optional[Dict]:
        pass
    @abstractmethod
    def set(self, key: str, value: Dict, expire: int = 3600):
        pass
    @abstractmethod
    def delete(self, key: str):
        pass
    @abstractmethod
    def acquire_lock(self, key: str, expire: int = 10) -> bool:
        pass
    @abstractmethod
    def release_lock(self, key: str):
        pass
class RedisIdempotentStorage(IdempotentStorage):
    """Redis存储实现"""
    def __init__(self, redis_client: redis.Redis):
        self.redis = redis_client
    def get(self, key: str) -> Optional[Dict]:
        data = self.redis.get(key)
        return json.loads(data) if data else None
    def set(self, key: str, value: Dict, expire: int = 3600):
        self.redis.setex(key, expire, json.dumps(value))
    def delete(self, key: str):
        self.redis.delete(key)
    def acquire_lock(self, key: str, expire: int = 10) -> bool:
        lock_key = f"{key}:lock"
        return bool(self.redis.setnx(lock_key, 1)) and bool(self.redis.expire(lock_key, expire))
    def release_lock(self, key: str):
        lock_key = f"{key}:lock"
        self.redis.delete(lock_key)
class MemoryIdempotentStorage(IdempotentStorage):
    """内存存储实现(适合单机测试)"""
    def __init__(self):
        self.storage = {}
        self.locks = {}
        self.expire_times = {}
        self._cleanup_thread = threading.Thread(target=self._cleanup_expired, daemon=True)
        self._cleanup_thread.start()
    def get(self, key: str) -> Optional[Dict]:
        if key in self.storage:
            if datetime.now() < self.expire_times.get(key, datetime.min):
                return self.storage[key]
            else:
                self.delete(key)
        return None
    def set(self, key: str, value: Dict, expire: int = 3600):
        self.storage[key] = value
        self.expire_times[key] = datetime.now() + timedelta(seconds=expire)
    def delete(self, key: str):
        self.storage.pop(key, None)
        self.expire_times.pop(key, None)
        self.locks.pop(key, None)
    def acquire_lock(self, key: str, expire: int = 10) -> bool:
        lock_key = f"{key}:lock"
        if lock_key not in self.locks:
            self.locks[lock_key] = datetime.now() + timedelta(seconds=expire)
            return True
        elif datetime.now() > self.locks[lock_key]:
            # 锁过期
            self.locks[lock_key] = datetime.now() + timedelta(seconds=expire)
            return True
        return False
    def release_lock(self, key: str):
        lock_key = f"{key}:lock"
        self.locks.pop(lock_key, None)
    def _cleanup_expired(self):
        while True:
            time.sleep(60)
            now = datetime.now()
            expired_keys = [k for k, v in self.expire_times.items() if now > v]
            for key in expired_keys:
                self.delete(key)
class IdempotentTool:
    """幂等工具主类"""
    def __init__(self, storage: IdempotentStorage, key_prefix: str = "idempotent"):
        self.storage = storage
        self.key_prefix = key_prefix
    def generate_key(self, func_name: str, args: tuple, kwargs: dict) -> str:
        """生成幂等key"""
        idempotent_key = kwargs.get('idempotent_key')
        if not idempotent_key:
            content = f"{func_name}:{json.dumps(args)}:{json.dumps(kwargs, sort_keys=True)}"
            idempotent_key = hashlib.md5(content.encode()).hexdigest()
        return f"{self.key_prefix}:{func_name}:{idempotent_key}"
    def idempotent(self, expire_time: int = 3600):
        """幂等装饰器"""
        def decorator(func):
            @wraps(func)
            def wrapper(*args, **kwargs):
                key = self.generate_key(func.__name__, args, kwargs)
                # 检查是否已处理
                existing_result = self.storage.get(key)
                if existing_result:
                    if existing_result['status'] == 'completed':
                        return existing_result['data']
                    elif existing_result['status'] == 'processing':
                        return {"error": "请求正在处理中"}, 409
                # 获取分布式锁
                if not self.storage.acquire_lock(key):
                    return {"error": "请求正在处理中"}, 409
                try:
                    # 标记为处理中
                    self.storage.set(key, {
                        'status': 'processing',
                        'timestamp': datetime.now().isoformat()
                    }, expire_time)
                    # 执行业务逻辑
                    result = func(*args, **kwargs)
                    # 更新为完成状态
                    self.storage.set(key, {
                        'status': 'completed',
                        'data': result,
                        'timestamp': datetime.now().isoformat()
                    }, expire_time)
                    return result
                except Exception as e:
                    # 发生异常,删除处理标记
                    self.storage.delete(key)
                    raise e
                finally:
                    # 释放锁
                    self.storage.release_lock(key)
            return wrapper
        return decorator
# 使用示例
redis_client = redis.Redis(host='localhost', port=6379, db=0)
storage = RedisIdempotentStorage(redis_client)
idempotent_tool = IdempotentTool(storage)
@idempotent_tool.idempotent(expire_time=7200)
def create_order(order_data: Dict) -> Dict:
    """创建订单(幂等)"""
    # 订单创建逻辑
    order_id = f"order_{datetime.now().timestamp()}"
    return {"order_id": order_id, "status": "created", "data": order_data}
# 调用示例
result = create_order({"product_id": "prod_123", "quantity": 2}, idempotent_key="unique_key_123")

使用建议

# 1. 选择合适的幂等方案
# - 简单应用:使用装饰器 + Redis
# - 复杂应用:使用分布式锁 + 数据库记录
# - 高并发场景:使用Redis + Lua脚本保证原子性
# 2. 幂等key生成策略
def generate_idempotent_key(request_obj):
    """生成具有业务意义的幂等key"""
    components = [
        request_obj.method,
        request_obj.path,
        request_obj.get_json().get('business_key', ''),
        request_obj.headers.get('X-Request-Id', '')
    ]
    return hashlib.md5('|'.join(components).encode()).hexdigest()
# 3. 幂等与事务结合
class TransactionIdempotentTool(IdempotentTool):
    """支持事务的幂等工具"""
    def idempotent_with_transaction(self, session_factory):
        """带事务的幂等装饰器"""
        def decorator(func):
            @wraps(func)
            def wrapper(*args, **kwargs):
                key = self.generate_key(func.__name__, args, kwargs)
                # 检查幂等性
                existing = self.storage.get(key)
                if existing and existing['status'] == 'completed':
                    return existing['data']
                # 使用数据库事务
                session = session_factory()
                try:
                    result = func(*args, **kwargs)
                    session.commit()
                    # 记录幂等结果
                    self.storage.set(key, {
                        'status': 'completed',
                        'data': result
                    })
                    return result
                except:
                    session.rollback()
                    raise
                finally:
                    session.close()
            return wrapper
        return decorator

关键点:

  1. 幂等key生成:基于请求内容或业务标识
  2. 存储方式:Redis、数据库、内存
  3. 并发控制:分布式锁防止重复处理
  4. 过期机制:防止内存泄漏和key堆积
  5. 错误处理:异常时清除幂等标记

选择方案时考虑:

  • 系统的并发量
  • 业务的重要性
  • 可用性和一致性要求
  • 运维复杂度

抱歉,评论功能暂时关闭!