fix(fam-edge): Gemini 多 key 从未真正轮转 - 流量全压在 key1

analyze_video/_generate_text 原来每次都从 api_keys[0] 开始遍历,只有当
key 的全部模型都失败才换下一个 key;由于 flash-lite 兜底基本总能在
key1 下成功,key2/3/4 实际零流量,配置的多 key 配额分散形同虚设。

新增 _rotated_keys():每次调用前记住上次轮转到的起点,按起点滚动排序
返回 key 列表并把起点推进到下一个,多次调用后流量均匀摊到全部 key。
model_calls 统计的 model 字段现在带上 ·key{n} 后缀,可以看出具体是哪个
key 在跑(ModelStats 页面对应拆出 Key 列展示)。

新增 4 个单测覆盖起点归零/逐次推进/回绕/单 key 情形。
This commit is contained in:
ericwyuan
2026-08-22 07:19:23 +08:00
parent 6b149a0dd0
commit 1ad6af3d5d
2 changed files with 95 additions and 37 deletions

View File

@@ -9,11 +9,23 @@ provider_name = "gemini"
整视频分析: 用 Files API 上传完整视频 -> generateContent 直出结构化 JSON
本地不切片、不抽帧Gemini 原生支持长视频)
多 Key 轮换2026-08-21 新增): 不同 Google Cloud 项目的 API Key 各自独立计费/配额,
config 的 `api_key` 为主 Key`extra_api_keys` 可以再配多个(各自项目的 Key
outer loop 按 Key 顺序尝试inner loop 才是原有的模型 fallback 链——因为 Gemini
Files API 上传的文件只能被同一个 Key/项目引用,换 Key 必须重新上传,所以每个 Key
都要重走一遍"上传 -> 模型链尝试 -> 删除",不是简单地在同一次上传后换 key 调用
多 Key 轮换2026-08-21 新增2026-08-22 改为真正均摊负载: 不同 Google Cloud
项目的 API Key 各自独立计费/配额,config 的 `api_key` 为主 Key`extra_api_keys`
可以再配多个(各自项目的 Keyouter loop 按 Key 顺序尝试inner loop 才是原有
的模型 fallback 链——因为 Gemini Files API 上传的文件只能被同一个 Key/项目引用,
换 Key 必须重新上传,所以每个 Key 都要重走一遍"上传 -> 模型链尝试 -> 删除"
实测发现的问题: 原实现每次都从 api_keys[0] 开始试,只有 0 号 key 的所有模型全部
失败才会换下一个 key但 flash-lite 兜底通常最终能成功,导致 0 号 key 几乎揽下
全部流量,其余 3 个 key 常年闲置——完全没有起到分摊配额的作用。现在改为
_rotated_keys():每次 analyze_video()/chat() 调用都从"上一次的下一个 key"开始
试起,调用结束(无论成败)就把起点往后挪一位,多次调用下来自然把请求均匀摊到
所有配置的 key 上,而不是"谁在前面谁扛所有流量"
每次实际使用的 key 会以 "key{N}"N 从 1 开始,对应 api_keys 里的原始下标)的
形式拼进 model_calls.model 字段(如 "gemini-flash-latest·key2"),这样现有的
按 provider+model 分组统计fam-core ui_api.py 的 /api/ui/model-stats不用改
schema 就能天然按 key 拆开显示,不需要新增字段/新迁移。
"""
import os
import time
@@ -50,6 +62,7 @@ class GeminiAdapter(BaseModelAdapter):
seen.add(resolved)
self.api_keys.append(resolved)
self.api_key = self.api_keys[0] if self.api_keys else '' # 向后兼容单 key 用法
self._key_rotation_idx = 0 # 下一次调用从哪个 key 起手(轮转,均摊负载用)
self.timeout = config.get('timeout', 600)
# 模型级独立超时(最终值,不参与编排层 ×N 放大): {model_name: seconds}
# 例: {"gemini-flash-lite-latest": 90}(按实测耗时 ×4 配置)
@@ -68,6 +81,17 @@ class GeminiAdapter(BaseModelAdapter):
return os.environ.get(raw[2:-1], '')
return raw
def _rotated_keys(self):
"""按当前轮转起点排序的 (原始下标从0开始, key) 列表;每调一次就把起点挪
到下一个 key多次调用下来把流量均匀摊到全部配置的 key 上。"""
n = len(self.api_keys)
if n == 0:
return []
start = self._key_rotation_idx % n
order = list(range(start, n)) + list(range(0, start))
self._key_rotation_idx = (start + 1) % n
return [(i, self.api_keys[i]) for i in order]
def health_check(self) -> bool:
if not self.api_key:
logger.warning("Gemini API Key 未配置,健康检查失败")
@@ -108,41 +132,42 @@ class GeminiAdapter(BaseModelAdapter):
prompt = self._build_video_prompt(known_members_context, event_start_time)
last_err = "no_key_available"
for idx, key in enumerate(self.api_keys):
for idx, key in self._rotated_keys():
key_label = f"key{idx + 1}"
file_uri, file_name = self._upload_file(video_path, key)
if not file_uri:
self._delete_file(file_name, key) # 即使未等到 ACTIVE也尽力清理
last_err = f"key{idx}_upload_failed"
last_err = f"{key_label}_upload_failed"
continue
try:
# 3 秒密度 + person_appearances 特征使输出 JSON 较大max_tokens 需足够高防截断
text = self._generate_video(file_uri, prompt, key,
text = self._generate_video(file_uri, prompt, key, key_label,
max_tokens=16384, temperature=0.2)
if text is None:
last_err = f"key{idx}_all_models_failed"
last_err = f"{key_label}_all_models_failed"
continue
try:
result = parse_vlm_json(text)
result = self._normalize(result)
if not result or 'events' not in result:
logger.error(f"Gemini 视频输出缺少 events: {text[:150]}")
last_err = f"key{idx}_missing_events"
last_err = f"{key_label}_missing_events"
continue
result['compute_provider'] = 'gemini'
self._cb.record_success()
logger.info(f"Gemini 整视频分析完成(key[{idx}]"
logger.info(f"Gemini 整视频分析完成({key_label}"
f"events={len(result.get('events', []))}")
return result
except VLMOutputInvalidError as e:
logger.error(f"Gemini 视频输出无法解析为 JSON: {e}")
last_err = f"key{idx}_invalid_json"
last_err = f"{key_label}_invalid_json"
continue
except requests.Timeout:
logger.warning(f"Gemini key[{idx}] 视频分析超时 ({self.timeout}s)")
last_err = f"key{idx}_timeout"
logger.warning(f"Gemini {key_label} 视频分析超时 ({self.timeout}s)")
last_err = f"{key_label}_timeout"
except Exception as e:
logger.error(f"Gemini key[{idx}] 视频分析异常: {e}")
last_err = f"key{idx}_exception"
logger.error(f"Gemini {key_label} 视频分析异常: {e}")
last_err = f"{key_label}_exception"
finally:
self._delete_file(file_name, key)
self._cb.record_failure()
@@ -250,7 +275,7 @@ class GeminiAdapter(BaseModelAdapter):
except Exception:
pass
def _generate_video(self, file_uri: str, prompt: str, api_key: str,
def _generate_video(self, file_uri: str, prompt: str, api_key: str, key_label: str,
max_tokens: int, temperature: float) -> Optional[str]:
"""带模型 fallback 链的 generateContent视频文件引用调用。"""
parts = [
@@ -258,11 +283,12 @@ class GeminiAdapter(BaseModelAdapter):
{"text": prompt},
]
for model in self.model_chain:
stat_model = f"{model}·{key_label}" # 供 model_calls 按 key 拆分统计
# 模型级独立超时优先;未单独配置的用适配器默认(可能已被编排层 ×N 放大)
model_timeout = self.model_timeouts.get(model, self.timeout)
for attempt in range(2):
if attempt == 0:
logger.info(f"Gemini [{model}] 本轮请求超时 {model_timeout}s")
logger.info(f"Gemini [{stat_model}] 本轮请求超时 {model_timeout}s")
started = datetime.now(timezone(timedelta(hours=8))).strftime('%Y-%m-%d %H:%M:%S')
t0 = time.time()
try:
@@ -276,14 +302,14 @@ class GeminiAdapter(BaseModelAdapter):
)
except requests.Timeout:
duration = time.time() - t0
self._emit_model_call(model, started, duration, False,
self._emit_model_call(stat_model, started, duration, False,
f"timeout({model_timeout}s)")
logger.warning(f"Gemini [{model}] 视频请求超时 ({model_timeout}s)")
logger.warning(f"Gemini [{stat_model}] 视频请求超时 ({model_timeout}s)")
break
except Exception as e:
duration = time.time() - t0
self._emit_model_call(model, started, duration, False, str(e))
logger.error(f"Gemini [{model}] 视频请求异常: {e}")
self._emit_model_call(stat_model, started, duration, False, str(e))
logger.error(f"Gemini [{stat_model}] 视频请求异常: {e}")
break
duration = time.time() - t0
@@ -294,29 +320,29 @@ class GeminiAdapter(BaseModelAdapter):
for p in (cands[0].get('content', {}) if cands else {}).get('parts', [])
).strip() if cands else ''
if text:
self._emit_model_call(model, started, duration, True)
self._emit_model_call(stat_model, started, duration, True)
if model != self.model_name:
logger.info(f"Gemini 主模型不可用,由 fallback [{model}] 出结果")
return text
self._emit_model_call(model, started, duration, False, "empty_text")
logger.warning(f"Gemini [{model}] 返回空文本")
self._emit_model_call(stat_model, started, duration, False, "empty_text")
logger.warning(f"Gemini [{stat_model}] 返回空文本")
continue
detail = resp.text[:150].replace('\n', ' ')
if resp.status_code == 429:
self._emit_model_call(model, started, duration, False, "429_quota")
logger.warning(f"Gemini [{model}] 429 配额耗尽,切换下一模型")
self._emit_model_call(stat_model, started, duration, False, "429_quota")
logger.warning(f"Gemini [{stat_model}] 429 配额耗尽,切换下一模型")
break
if resp.status_code == 503:
if attempt == 0:
self._emit_model_call(model, started, duration, False, "503_overload_retry")
logger.warning(f"Gemini [{model}] 503 过载3s 后重试")
self._emit_model_call(stat_model, started, duration, False, "503_overload_retry")
logger.warning(f"Gemini [{stat_model}] 503 过载3s 后重试")
time.sleep(3)
continue
self._emit_model_call(model, started, duration, False, "503_overload")
self._emit_model_call(stat_model, started, duration, False, "503_overload")
break
self._emit_model_call(model, started, duration, False,
self._emit_model_call(stat_model, started, duration, False,
f"http_{resp.status_code}")
logger.warning(f"Gemini [{model}] HTTP {resp.status_code}: {detail}")
logger.warning(f"Gemini [{stat_model}] HTTP {resp.status_code}: {detail}")
break
return None
@@ -383,8 +409,10 @@ class GeminiAdapter(BaseModelAdapter):
return result
def _generate_text(self, text: str, max_tokens: int, temperature: float) -> Optional[str]:
"""纯文本 generateContent按 key 轮换 × 模型 fallback 链依次尝试。"""
for idx, api_key in enumerate(self.api_keys):
"""纯文本 generateContent按 key 轮换(同 analyze_video 共用一套轮转起点)
× 模型 fallback 链依次尝试。"""
for idx, api_key in self._rotated_keys():
key_label = f"key{idx + 1}"
for model in self.model_chain:
try:
resp = requests.post(
@@ -396,10 +424,10 @@ class GeminiAdapter(BaseModelAdapter):
timeout=self.timeout
)
except requests.Timeout:
logger.warning(f"Gemini key[{idx}] [{model}] 问答超时")
logger.warning(f"Gemini {key_label} [{model}] 问答超时")
continue
except Exception as e:
logger.error(f"Gemini key[{idx}] [{model}] 问答异常: {e}")
logger.error(f"Gemini {key_label} [{model}] 问答异常: {e}")
continue
if resp.status_code == 200:
cands = resp.json().get('candidates', [])
@@ -410,7 +438,7 @@ class GeminiAdapter(BaseModelAdapter):
if out:
return out
elif resp.status_code == 429:
logger.warning(f"Gemini key[{idx}] [{model}] 429切换下一模型/Key")
logger.warning(f"Gemini {key_label} [{model}] 429切换下一模型/Key")
continue
return None