You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python异步任务场景下带限流的OpenAI API稳健封装技术问询

Django异步任务中OpenAI API封装与限流优化问题

我正在开发Django项目,需为异步任务构建稳健高效的OpenAI API封装,核心目标包括:

  • API封装设计:创建支持模型选择与提示生成的封装类
  • 限流实现:处理每分钟请求数(RPM)与每分钟令牌数(TPM)双重限制
  • 异步任务管理:通过Celery等框架管控API调用,避免触发限流

当前构思

  • 封装类:负责OpenAI API连接、模型选择与提示发送
  • 提示生成:专用类/方法构建管理提示结构
  • 限流机制:结合Celery实现速率管控,但当前Redis TokenBucket实现存在并发问题(Celery Worker并发数为1仍触发其他API的"Too Many Requests")

咨询问题

  1. 最佳实践:构建此类封装的推荐模式,尤其是错误处理与效率优化
  2. 限流技术:同时支持RPM与TPM的高效方案,及请求前预估令牌消耗的方法
  3. 异步集成:如何集成Celery或其他框架,在不违反限流规则下管理任务队列

附上我的TokenBucket类与Celery刷新任务代码

class TokenBucketManager:
    def __init__(self, endpoint, tokens, ttl, high_frequency=False):
        self.redis = caches["token_cache"]
        self.endpoint = endpoint
        self.tokens = tokens
        self.ttl = ttl
        self.high_frequency = high_frequency
        
    def _bucket_key(self):
        return f"token_bucket:{self.endpoint}"

    def _keepalive_key(self):
        return f"token_bucket_keepalive:{self.endpoint}"

    def _sync_key(self):
        return f"token_bucket_sync:{self.endpoint}"

    def token_bucket_exists(self):
        bucket_key = self._bucket_key() 
        return self.redis.has_key(bucket_key)

    def keep_token_alive(self):
        keep_alive_key = self._keepalive_key()
        token_bucket_key = self._bucket_key()
        sync_key = self._sync_key()
        task_token_bucket.apply_async(
            countdown=self.ttl - 1, 
            kwargs={"keep_alive_key": keep_alive_key, "token_bucket_key": token_bucket_key, "sync_key":sync_key, "ttl": self.ttl, "tokens": self.tokens, "is_sync": False}, 
            queue="token_handler"
            )
        return keep_alive_key
        
    def set_token_bucket(self):
        key = self._bucket_key()
        with self.redis.lock(key):
            if not self.redis.has_key(key):
                self.redis.set(key, self.tokens, self.ttl)
                self.keep_token_alive() 
                logger.info(f"Setting token bucket for {key} with {self.tokens} tokens and ttl {self.ttl}")
        
    def consume_if_possible(self, amount=1):
        
            if not self.token_bucket_exists():
                self.set_token_bucket()
                
            token_bucket_key = self._bucket_key()
            keep_alive_key = self._keepalive_key()
            self.redis.set(keep_alive_key, 1, self.ttl * 2)
            if self.redis.has_key(token_bucket_key):
                with self.redis.lock(token_bucket_key):
                    tokens_left = self.redis.get(token_bucket_key)
                    if tokens_left is not None and tokens_left > amount:
                        try:
                            self.redis.decr(token_bucket_key, amount)
                        except:
                            return False 
                        return True
                    else:
                        return False
            else:
                return False 


    def get_tokens_ttl(self):
        ttl = self.redis.ttl(self._bucket_key())
        return ttl if ttl > 0 else self.ttl


    def sync_bucket_to_external_api(self, retry_after):
        if self.high_frequency:
            pass 
        logger.info(f"Syncing token bucket for {self._bucket_key()} to external API, Retry After: {retry_after}")
        token_bucket_key = self._bucket_key()
        sync_key = self._sync_key()
        keepalive_key = self._keepalive_key()
        with self.redis.lock(token_bucket_key):
            self.redis.set(token_bucket_key, 0, retry_after)
        self.redis.set(sync_key, 1, self.ttl * 2)
        task_token_bucket.apply_async(
            countdown=retry_after,
            kwargs={"keep_alive_key": keepalive_key, "token_bucket_key": token_bucket_key, "sync_key":sync_key, "ttl": self.ttl, "tokens": self.tokens, "is_sync": True},
            queue="token_handler")

用于按需刷新令牌桶的Celery任务

@shared_task(bind=True)
def task_token_bucket(self, keep_alive_key, token_bucket_key, 
sync_key, ttl, tokens, is_sync=False):
    redis = caches["token_cache"]

    if redis.has_key(sync_key):
        if not is_sync:
            return 
        else:
            with redis.lock(sync_key):
                if redis.has_key(sync_key):
                    redis.delete(sync_key)
                else:
                    return
            
    if redis.has_key(keep_alive_key):
        with redis.lock(token_bucket_key, 30):
            if redis.has_key(token_bucket_key):
                redis.set(token_bucket_key, tokens, ttl)
                redis.delete(keep_alive_key)
                logger.info(f"Token bucket for {token_bucket_key} refreshed.")
                task_token_bucket.apply_async(
                    countdown=ttl,
                    kwargs={"keep_alive_key": keep_alive_key, "sync_key": sync_key, "token_bucket_key": token_bucket_key, "ttl": ttl, "tokens": tokens},
                    queue="token_handler"
                )
                logger.info(f"Renewal task for token bucket {token_bucket_key} scheduled.")
            else:
                logger.info(f"Token bucket for {token_bucket_key} not refreshed.")
    else:
        logger.info(f"Token bucket for {token_bucket_key} not refreshed.")

内容的提问来源于stack exchange,提问作者SeRo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 03:01:03