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")
咨询问题
- 最佳实践:构建此类封装的推荐模式,尤其是错误处理与效率优化
- 限流技术:同时支持RPM与TPM的高效方案,及请求前预估令牌消耗的方法
- 异步集成:如何集成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
相关产品推荐
相关产品推荐

