本文目录导读:

我来为你提供几个漏桶算法的具体案例,从简单到复杂,帮助理解其应用。
基本漏桶算法实现(Python)
import time
import threading
from collections import deque
class LeakyBucket:
"""漏桶算法基本实现"""
def __init__(self, capacity, leak_rate):
"""
capacity: 桶的容量(最大请求数)
leak_rate: 漏出速率(每秒处理的请求数)
"""
self.capacity = capacity
self.leak_rate = leak_rate
self.water = 0 # 当前桶中的水量(请求数)
self.last_time = time.time() # 上次漏水时间
self.lock = threading.Lock()
def allow_request(self):
"""判断是否允许请求通过"""
with self.lock:
# 先漏水
current_time = time.time()
elapsed = current_time - self.last_time
# 计算这段时间内漏出的水量
leaked = elapsed * self.leak_rate
self.water = max(0, self.water - leaked)
self.last_time = current_time
# 检查桶是否已满
if self.water < self.capacity:
self.water += 1
return True
else:
return False
# 使用示例
def test_basic_leaky_bucket():
# 创建一个容量为5,漏出速率为2个/秒的桶
bucket = LeakyBucket(capacity=5, leak_rate=2)
# 模拟连续请求
for i in range(10):
allowed = bucket.allow_request()
print(f"请求 {i+1}: {'允许' if allowed else '拒绝'}")
time.sleep(0.1) # 每0.1秒发送一个请求
# 等待漏完
print("\n等待2秒后...")
time.sleep(2)
print(f"当前水量: {bucket.water}")
电商平台限流案例
import time
import threading
from datetime import datetime
class EcommerceRateLimiter:
"""电商平台下单限流器"""
def __init__(self, order_capacity=100, leak_rate=10):
"""
order_capacity: 订单桶容量(同时处理的最大订单数)
leak_rate: 处理速率(每秒处理的订单数)
"""
self.capacity = order_capacity
self.leak_rate = leak_rate
self.water = 0
self.last_time = time.time()
self.request_log = []
self.lock = threading.Lock()
def process_order(self, user_id, order_data):
"""处理订单"""
with self.lock:
self._leak_water()
if self.water >= self.capacity:
return {
'success': False,
'message': '系统繁忙,请稍后重试',
'code': 429
}
self.water += 1
log_entry = {
'user_id': user_id,
'timestamp': datetime.now(),
'order_data': order_data
}
self.request_log.append(log_entry)
# 模拟订单处理耗时
processing_time = 1 / self.leak_rate
threading.Timer(processing_time, self._complete_order,
args=(log_entry,)).start()
return {
'success': True,
'message': '订单已接受处理',
'code': 200,
'order_id': len(self.request_log)
}
def _leak_water(self):
"""漏水处理"""
current_time = time.time()
elapsed = current_time - self.last_time
leaked = elapsed * self.leak_rate
self.water = max(0, self.water - leaked)
self.last_time = current_time
def _complete_order(self, log_entry):
"""完成订单处理"""
with self.lock:
self.water = max(0, self.water - 1)
print(f"订单 {log_entry['order_id']} 处理完成")
def get_status(self):
"""获取当前状态"""
self._leak_water()
return {
'current_load': self.water,
'capacity': self.capacity,
'utilization': self.water / self.capacity,
'leak_rate': self.leak_rate
}
# 使用示例
def test_ecommerce_limiter():
limiter = EcommerceRateLimiter(order_capacity=5, leak_rate=1)
# 模拟并发用户下单
print("模拟10个用户同时下单:")
for user_id in range(10):
result = limiter.process_order(f"user_{user_id}", {"product": "手机"})
print(f"用户 {user_id}: {result['message']}")
print(f"\n系统状态: {limiter.get_status()}")
# 等待处理完成
time.sleep(3)
print(f"\n3秒后系统状态: {limiter.get_status()}")
API网关漏桶限流
import time
import hashlib
from collections import defaultdict
class APIGatewayRateLimiter:
"""API网关多用户漏桶限流器"""
def __init__(self, global_capacity=1000, global_leak_rate=100,
user_capacity=10, user_leak_rate=5):
self.global_bucket = LeakyBucket(global_capacity, global_leak_rate)
self.user_buckets = defaultdict(
lambda: LeakyBucket(user_capacity, user_leak_rate)
)
self.request_history = []
def api_request(self, user_id, endpoint):
"""处理API请求"""
request_id = self._generate_request_id(user_id, endpoint)
# 双重限制:全局和用户级别
if not self.global_bucket.allow_request():
return {
'status': 429,
'message': '全局请求过多',
'request_id': request_id
}
if not self.user_buckets[user_id].allow_request():
return {
'status': 429,
'message': '个人请求频率超过限制',
'request_id': request_id
}
# 记录请求
request_info = {
'request_id': request_id,
'user_id': user_id,
'endpoint': endpoint,
'timestamp': time.time()
}
self.request_history.append(request_info)
return {
'status': 200,
'message': '请求成功',
'request_id': request_id,
'data': request_info
}
def _generate_request_id(self, user_id, endpoint):
"""生成请求ID"""
raw_id = f"{user_id}:{endpoint}:{time.time()}"
return hashlib.md5(raw_id.encode()).hexdigest()[:8]
def get_user_requests(self, user_id, time_window=60):
"""获取用户在指定时间窗口内的请求数"""
current_time = time.time()
user_requests = [
r for r in self.request_history
if r['user_id'] == user_id
and current_time - r['timestamp'] <= time_window
]
return len(user_requests)
# 使用示例
def test_api_gateway():
limiter = APIGatewayRateLimiter()
print("模拟用户频繁调用API:")
for i in range(15):
result = limiter.api_request("user_123", "/api/data")
print(f"请求 {i+1}: HTTP {result['status']} - {result['message']}")
time.sleep(0.2)
print(f"\n用户1分钟内请求数: {limiter.get_user_requests('user_123')}")
消息队列漏桶应用
import time
import queue
import threading
from dataclasses import dataclass
@dataclass
class Message:
"""消息数据类"""
id: int
content: str
priority: int = 1
class MessageQueueWithLeakyBucket:
"""基于漏桶的消息队列限流器"""
def __init__(self, bucket_capacity=50, process_rate=10):
"""
bucket_capacity: 桶容量
process_rate: 处理速率(每秒处理消息数)
"""
self.bucket_capacity = bucket_capacity
self.process_rate = process_rate
self.message_queue = queue.Queue()
self.water = 0
self.last_time = time.time()
self.processing = False
self.lock = threading.Lock()
self.processed_count = 0
def add_message(self, content):
"""添加消息到队列"""
with self.lock:
# 检查桶容量
if self.water >= self.bucket_capacity:
print(f"队列已满,拒绝消息: {content}")
return False
message = Message(
id=time.time_ns(),
content=content
)
self.message_queue.put(message)
self.water += 1
# 启动处理线程(如果未在运行)
if not self.processing:
self._start_processing()
return True
def _start_processing(self):
"""启动处理进程"""
self.processing = True
processing_thread = threading.Thread(target=self._process_messages)
processing_thread.daemon = True
processing_thread.start()
def _process_messages(self):
"""处理消息"""
while self.processing:
with self.lock:
# 漏水
current_time = time.time()
elapsed = current_time - self.last_time
leaked = elapsed * self.process_rate
self.water = max(0, self.water - leaked)
self.last_time = current_time
try:
# 获取消息(非阻塞)
message = self.message_queue.get_nowait()
print(f"处理消息: {message.content}")
self.processed_count += 1
self.message_queue.task_done()
# 控制处理速率
time.sleep(1 / self.process_rate)
except queue.Empty:
if self.message_queue.empty():
self.processing = False
def get_queue_status(self):
"""获取队列状态"""
with self.lock:
current_time = time.time()
elapsed = current_time - self.last_time
leaked = elapsed * self.process_rate
current_water = max(0, self.water - leaked)
return {
'current_load': current_water,
'queue_size': self.message_queue.qsize(),
'processed_count': self.processed_count,
'bucket_capacity': self.bucket_capacity
}
# 使用示例
def test_message_queue():
mq = MessageQueueWithLeakyBucket(bucket_capacity=5, process_rate=2)
print("发送8条消息:")
for i in range(8):
success = mq.add_message(f"消息-{i+1}")
print(f"添加消息 {i+1}: {'成功' if success else '失败'}")
time.sleep(0.1)
time.sleep(2)
print(f"\n队列状态: {mq.get_queue_status()}")
实际应用场景总结
class LeakyBucketDemo:
"""漏桶算法完整演示"""
@staticmethod
def demo_pipeline_processing():
"""管道流水线处理演示"""
print("="*50)
print("漏桶算法模拟管道流水线")
print("="*50)
# 模拟生产速率与消费速率
bucket = LeakyBucket(capacity=3, leak_rate=1) # 每秒处理1个
# 模拟生产过程
for i in range(6):
allowed = bucket.allow_request()
status = "通过" if allowed else "拒绝"
print(f"时间{i+1}: 事件{status}")
time.sleep(0.5) # 每0.5秒生产一个
@staticmethod
def demo_burst_handling():
"""突发流量处理演示"""
print("="*50)
print("突发流量处理")
print("="*50)
bucket = LeakyBucket(capacity=10, leak_rate=2)
# 模拟突发流量
burst_requests = 10
for i in range(burst_requests):
allowed = bucket.allow_request()
print(f"突发请求 {i+1}: {'接受' if allowed else '拒绝'}")
print("\n等待处理...")
time.sleep(2)
print(f"2秒后还能接受 {bucket.capacity - bucket.water} 个请求")
# 运行演示
if __name__ == "__main__":
demo = LeakyBucketDemo()
# 运行各个演示
print("\n1. 基本漏桶算法")
test_basic_leaky_bucket()
print("\n2. 电商平台限流")
test_ecommerce_limiter()
print("\n3. API网关限流")
test_api_gateway()
print("\n4. 消息队列应用")
test_message_queue()
print("\n5. 场景演示")
demo.demo_pipeline_processing()
demo.demo_burst_handling()
- 平滑处理:漏桶将突发请求平滑处理,输出速率恒定
- 流量整形:适合需要稳定速率的场景
- 内存保护:通过桶容量限制内存占用
- 实现简单:逻辑清晰,易于实现和维护
这些案例展示了漏桶算法在实时流量控制、API保护、消息处理等场景的应用,你可以根据实际需求选择适合的实现方式。