Python熔断工具案例如何封装服务熔断

wen python案例 33

本文目录导读:

Python熔断工具案例如何封装服务熔断

  1. 基础熔断器封装
  2. 装饰器方式的熔断器
  3. 实际应用示例
  4. 使用示例
  5. 最佳实践建议

我来详细介绍Python服务熔断的封装方法,以常用的pybreaker库为例。

基础熔断器封装

import pybreaker
from functools import wraps
import time
import logging
from typing import Callable, Any, Optional
logger = logging.getLogger(__name__)
class CircuitBreaker:
    """服务熔断器封装类"""
    def __init__(self, 
                 name: str = "default",
                 fail_max: int = 5,
                 reset_timeout: int = 60,
                 exclude_exceptions: tuple = None):
        """
        初始化熔断器
        Args:
            name: 熔断器名称
            fail_max: 最大失败次数
            reset_timeout: 重置超时时间(秒)
            exclude_exceptions: 不触发熔断的异常类型
        """
        self.name = name
        self.breaker = pybreaker.CircuitBreaker(
            fail_max=fail_max,
            reset_timeout=reset_timeout,
            exclude=exclude_exceptions or ()
        )
        self.state = "closed"  # 熔断器状态
        self.last_failure_time = None
        self.failure_count = 0
    def call(self, func: Callable, *args, **kwargs) -> Any:
        """执行被保护的方法调用"""
        try:
            result = self.breaker.call(func, *args, **kwargs)
            self._update_state("closed")
            return result
        except pybreaker.CircuitBreakerError:
            self._update_state("open")
            raise ServiceUnavailableError(f"Service {self.name} is unavailable")
        except Exception as e:
            self._update_state("half-open" if self.breaker.state == pybreaker.STATE_HALF_OPEN else "closed")
            raise
    def _update_state(self, state: str):
        """更新熔断器状态"""
        old_state = self.state
        self.state = state
        if old_state != state:
            logger.info(f"Circuit breaker {self.name} state changed: {old_state} -> {state}")
    def get_state(self) -> dict:
        """获取熔断器状态信息"""
        return {
            "name": self.name,
            "state": self.state,
            "failure_count": self.breaker.fail_counter,
            "reset_timeout": self.breaker.reset_timeout
        }
class ServiceUnavailableError(Exception):
    """服务不可用异常"""
    pass

装饰器方式的熔断器

class CircuitBreakerDecorator:
    """装饰器方式的熔断器"""
    _breakers = {}  # 存储所有熔断器实例
    @classmethod
    def breaker(cls, 
                name: str = None,
                fail_max: int = 3,
                reset_timeout: int = 30,
                fallback_func: Callable = None):
        """
        熔断器装饰器
        Args:
            name: 熔断器名称
            fail_max: 最大失败次数
            reset_timeout: 重置超时时间
            fallback_func: 降级函数
        """
        def decorator(func):
            breaker_name = name or func.__name__
            if breaker_name not in cls._breakers:
                cls._breakers[breaker_name] = CircuitBreaker(
                    name=breaker_name,
                    fail_max=fail_max,
                    reset_timeout=reset_timeout
                )
            @wraps(func)
            def wrapper(*args, **kwargs):
                circuit_breaker = cls._breakers[breaker_name]
                try:
                    return circuit_breaker.call(func, *args, **kwargs)
                except ServiceUnavailableError as e:
                    logger.warning(f"Circuit breaker {breaker_name} is open, using fallback")
                    if fallback_func:
                        return fallback_func(*args, **kwargs)
                    raise
            return wrapper
        return decorator

实际应用示例

import requests
from typing import Dict, List
# 服务接口封装
class UserServiceClient:
    """用户服务客户端"""
    def __init__(self, base_url: str):
        self.base_url = base_url
        self.session = requests.Session()
    @CircuitBreakerDecorator.breaker(
        name="user_service",
        fail_max=3,
        reset_timeout=30,
        fallback_func=lambda *args: {"error": "Service unavailable", "data": None}
    )
    def get_user(self, user_id: int) -> Dict:
        """获取用户信息"""
        response = self.session.get(
            f"{self.base_url}/users/{user_id}",
            timeout=5
        )
        response.raise_for_status()
        return response.json()
    @CircuitBreakerDecorator.breaker(
        name="user_list_service",
        fail_max=5,
        reset_timeout=60
    )
    def get_users(self, page: int = 1) -> List[Dict]:
        """获取用户列表"""
        response = self.session.get(
            f"{self.base_url}/users",
            params={"page": page},
            timeout=10
        )
        response.raise_for_status()
        return response.json()
# 高级熔断器实现
class AdvancedCircuitBreaker:
    """高级熔断器,支持滑动窗口和半开状态"""
    def __init__(self, 
                 name: str,
                 window_size: int = 10,
                 failure_threshold: float = 0.5,
                 half_open_max_calls: int = 3,
                 reset_timeout: int = 30):
        self.name = name
        self.window_size = window_size
        self.failure_threshold = failure_threshold
        self.half_open_max_calls = half_open_max_calls
        self.reset_timeout = reset_timeout
        self.state = "CLOSED"
        self.requests = []  # 滑动窗口中的请求记录
        self.last_failure_time = None
        self.half_open_calls = 0
    def _update_sliding_window(self, success: bool):
        """更新滑动窗口"""
        current_time = time.time()
        self.requests.append({
            "timestamp": current_time,
            "success": success
        })
        # 移除超过窗口时间的记录
        self.requests = [
            r for r in self.requests 
            if current_time - r["timestamp"] <= self.window_size
        ]
    def _get_failure_rate(self) -> float:
        """计算失败率"""
        if not self.requests:
            return 0.0
        failures = sum(1 for r in self.requests if not r["success"])
        return failures / len(self.requests)
    def call(self, func: Callable, *args, **kwargs) -> Any:
        """调用受保护的方法"""
        if self.state == "OPEN":
            if time.time() - self.last_failure_time >= self.reset_timeout:
                self.state = "HALF_OPEN"
                self.half_open_calls = 0
                logger.info(f"Circuit breaker {self.name} changed to HALF_OPEN")
            else:
                raise ServiceUnavailableError(f"Service {self.name} is unavailable")
        if self.state == "HALF_OPEN":
            if self.half_open_calls >= self.half_open_max_calls:
                raise ServiceUnavailableError(f"Too many half-open calls for {self.name}")
            self.half_open_calls += 1
        try:
            result = func(*args, **kwargs)
            self._update_sliding_window(success=True)
            if self.state == "HALF_OPEN":
                self.state = "CLOSED"
                self.half_open_calls = 0
                logger.info(f"Circuit breaker {self.name} reset to CLOSED")
            return result
        except Exception as e:
            self._update_sliding_window(success=False)
            self.last_failure_time = time.time()
            if self.state == "HALF_OPEN":
                self.state = "OPEN"
                logger.info(f"Circuit breaker {self.name} changed to OPEN")
            elif self._get_failure_rate() >= self.failure_threshold:
                self.state = "OPEN"
                logger.info(f"Circuit breaker {self.name} opened due to high failure rate")
            raise e

使用示例

# 服务配置
SERVICE_CONFIG = {
    "user_service": {"url": "http://user-service:8080", "timeout": 5},
    "order_service": {"url": "http://order-service:8080", "timeout": 10}
}
# 创建服务客户端
user_client = UserServiceClient(SERVICE_CONFIG["user_service"]["url"])
# 使用熔断器
def test_circuit_breaker():
    """测试熔断器功能"""
    # 正常调用
    try:
        user = user_client.get_user(1)
        print(f"User: {user}")
    except ServiceUnavailableError:
        print("User service is unavailable now")
    except Exception as e:
        print(f"Other error: {e}")
    # 获取熔断器状态
    breaker_info = CircuitBreakerDecorator._breakers["user_service"].get_state()
    print(f"Breaker state: {breaker_info}")
    # 批量测试
    for i in range(10):
        try:
            user = user_client.get_user(i)
            print(f"Request {i}: Success")
        except ServiceUnavailableError:
            print(f"Request {i}: Circuit breaker open")
        except Exception as e:
            print(f"Request {i}: Error - {e}")
        time.sleep(0.5)
if __name__ == "__main__":
    # 配置日志
    logging.basicConfig(level=logging.INFO)
    # 运行测试
    test_circuit_breaker()

最佳实践建议

  1. 合理设置阈值:根据服务特点和SLA要求调整
  2. 实现降级策略:提供备选方案或缓存数据
  3. 监控告警:记录熔断器状态变化
  4. 异步支持:对于异步框架提供对应的实现
  5. 分布式场景:考虑使用Redis等共享存储

这个封装示例提供了完整的熔断机制,可以根据实际需求进行调整和扩展。

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