本文目录导读:

我来详细介绍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()
最佳实践建议
- 合理设置阈值:根据服务特点和SLA要求调整
- 实现降级策略:提供备选方案或缓存数据
- 监控告警:记录熔断器状态变化
- 异步支持:对于异步框架提供对应的实现
- 分布式场景:考虑使用Redis等共享存储
这个封装示例提供了完整的熔断机制,可以根据实际需求进行调整和扩展。