漏桶算法案例

wen java案例 1

本文目录导读:

漏桶算法案例

  1. 基本漏桶算法实现(Python)
  2. 电商平台限流案例
  3. API网关漏桶限流
  4. 消息队列漏桶应用
  5. 实际应用场景总结

我来为你提供几个漏桶算法的具体案例,从简单到复杂,帮助理解其应用。

基本漏桶算法实现(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()
  1. 平滑处理:漏桶将突发请求平滑处理,输出速率恒定
  2. 流量整形:适合需要稳定速率的场景
  3. 内存保护:通过桶容量限制内存占用
  4. 实现简单:逻辑清晰,易于实现和维护

这些案例展示了漏桶算法在实时流量控制、API保护、消息处理等场景的应用,你可以根据实际需求选择适合的实现方式。

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