fix(resm): openalex退避匹配放宽Insufficient budget文案, 自触发链加单链去重

- OpenAlex 429 文案已从 "Insufficient credits" 变为 "Insufficient budget"
  (每日额度, UTC 午夜重置), 原精确匹配失效落入 2 分钟短退避分支导致空转,
  放宽为匹配 "Insufficient" 统一走 1 小时退避。
- 新增 claim_chain: 同名自触发链只允许一条存活。链自续发携带 chain_id,
  无 id 的调用(beat/手动/历史任务)接管为新链, 旧链下一轮自动退出,
  已繁殖的多条并行副本部署后自动收敛到单链。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
caoqianming 2026-07-21 14:41:59 +08:00
parent 2fe9c4cd1f
commit adf34b4bde
1 changed files with 34 additions and 5 deletions

View File

@ -13,6 +13,7 @@ from datetime import datetime, timedelta
import random import random
from .pdf_utils import _is_elsevier_preview_pdf from .pdf_utils import _is_elsevier_preview_pdf
from .d_oaurl import download_from_url_playwright from .d_oaurl import download_from_url_playwright
from uuid import uuid4
import asyncio import asyncio
import sys import sys
import os import os
@ -492,6 +493,24 @@ def touch_alive(def_name: str):
"""标记该自触发链仍在运行。""" """标记该自触发链仍在运行。"""
cache.set(def_name + ":alive", 1, timeout=ALIVE_TTL) cache.set(def_name + ":alive", 1, timeout=ALIVE_TTL)
def claim_chain(def_name: str, chain_id: str):
"""自触发链去重: 同名任务同时只允许一条链存活。
cache 里记录当前唯一合法链的 chain_id chain_id 的调用(beat 点火/手动 .delay/
历史遗留任务)视为新链接管, 旧链下一轮发现 chain_id 不匹配即自杀, 收敛到单链
返回本任务应携带的 chain_id; 返回 None 表示本任务是重复链, 应立即退出
"""
key = def_name + ":chain"
current = cache.get(key)
if chain_id and current and chain_id != current:
return None # 已有更新的链在跑, 本链退出
if not chain_id:
chain_id = uuid4().hex # 新点火: 接管成为唯一链
cache.set(key, chain_id, timeout=None)
elif not current:
cache.set(key, chain_id, timeout=None) # cache 丢失(如 redis 重启)后自愈
return chain_id
def is_alive(def_name: str): def is_alive(def_name: str):
return cache.get(def_name + ":alive") is not None return cache.get(def_name + ":alive") is not None
@ -506,10 +525,13 @@ def ensure_fetch_running():
return f"ensure_fetch_running started: {started}" return f"ensure_fetch_running started: {started}"
@shared_task(base=CustomTask) @shared_task(base=CustomTask)
def get_pdf_from_openalex(number_of_task: int =10): def get_pdf_from_openalex(number_of_task: int =10, chain_id: str = None):
def_name = get_pdf_from_openalex.name def_name = get_pdf_from_openalex.name
if not show_task_run(def_name): if not show_task_run(def_name):
return "stoped" return "stoped"
chain_id = claim_chain(def_name, chain_id)
if chain_id is None:
return "duplicate chain, exit"
touch_alive(def_name) touch_alive(def_name)
# 限流退避中: 不打 API, 慢节奏自重发只为维持 alive, 等 exceed 标记自然过期。 # 限流退避中: 不打 API, 慢节奏自重发只为维持 alive, 等 exceed 标记自然过期。
@ -517,7 +539,7 @@ def get_pdf_from_openalex(number_of_task: int =10):
if cache.get("openalex_api_exceed"): if cache.get("openalex_api_exceed"):
current_app.send_task( current_app.send_task(
"apps.resm.tasks.get_pdf_from_openalex", "apps.resm.tasks.get_pdf_from_openalex",
kwargs={"number_of_task": number_of_task}, kwargs={"number_of_task": number_of_task, "chain_id": chain_id},
countdown=60, countdown=60,
) )
return "openalex_api_exceed, backing off" return "openalex_api_exceed, backing off"
@ -545,6 +567,7 @@ def get_pdf_from_openalex(number_of_task: int =10):
"apps.resm.tasks.get_pdf_from_openalex", "apps.resm.tasks.get_pdf_from_openalex",
kwargs={ kwargs={
"number_of_task": number_of_task, "number_of_task": number_of_task,
"chain_id": chain_id,
}, },
countdown=countdown, countdown=countdown,
) )
@ -637,13 +660,16 @@ def _elsevier_fetch_pdf(req, paper):
@shared_task(base=CustomTask) @shared_task(base=CustomTask)
def get_abstract_from_elsevier(number_of_task:int = 20, exclude_failed:bool=True, def get_abstract_from_elsevier(number_of_task:int = 20, exclude_failed:bool=True,
pdf_number_of_task:int = 20): pdf_number_of_task:int = 20, chain_id: str = None):
"""Elsevier 单端点合并任务: 同一 DOI 先取 XML(摘要/全文标记), 若有全文则内联 """Elsevier 单端点合并任务: 同一 DOI 先取 XML(摘要/全文标记), 若有全文则内联
再取一次 PDF; 并补抓历史上已有全文标记但缺 PDF 的论文 get_pdf_from_elsevier 已并入 再取一次 PDF; 并补抓历史上已有全文标记但缺 PDF 的论文 get_pdf_from_elsevier 已并入
number_of_task: 阶段1(摘要+内联 PDF)每轮上限; pdf_number_of_task: 阶段2(存量补 PDF)每轮上限""" number_of_task: 阶段1(摘要+内联 PDF)每轮上限; pdf_number_of_task: 阶段2(存量补 PDF)每轮上限"""
def_name = get_abstract_from_elsevier.name def_name = get_abstract_from_elsevier.name
if not show_task_run(def_name): if not show_task_run(def_name):
return "stoped" return "stoped"
chain_id = claim_chain(def_name, chain_id)
if chain_id is None:
return "duplicate chain, exit"
touch_alive(def_name) touch_alive(def_name)
# 待抓摘要(并顺带取 PDF) # 待抓摘要(并顺带取 PDF)
@ -713,6 +739,7 @@ def get_abstract_from_elsevier(number_of_task:int = 20, exclude_failed:bool=True
"number_of_task": number_of_task, "number_of_task": number_of_task,
"exclude_failed": exclude_failed, "exclude_failed": exclude_failed,
"pdf_number_of_task": pdf_number_of_task, "pdf_number_of_task": pdf_number_of_task,
"chain_id": chain_id,
}, },
countdown=countdown, countdown=countdown,
) )
@ -870,10 +897,12 @@ def save_pdf_from_openalex(paper:Paper):
message = res.json().get("message", "") message = res.json().get("message", "")
except ValueError: except ValueError:
message = res.text message = res.text
if "Insufficient credits" in message: # 文案历史上出现过 "Insufficient credits" 和 "Insufficient budget"(每日额度,
# UTC 午夜重置), 放宽到 "Insufficient" 统一匹配, 避免落进 2 分钟短退避分支空转
if "Insufficient" in message:
# 额度耗尽: 退避 1 小时 # 额度耗尽: 退避 1 小时
cache.set("openalex_api_exceed", True, timeout=3600) cache.set("openalex_api_exceed", True, timeout=3600)
return "openalex_pdf_error: Insufficient credits" return f"openalex_pdf_error: {message[:100]}"
# 普通限流(请求过频): 短退避 2 分钟, 避免立刻重试再撞 429 # 普通限流(请求过频): 短退避 2 分钟, 避免立刻重试再撞 429
cache.set("openalex_api_exceed", True, timeout=120) cache.set("openalex_api_exceed", True, timeout=120)
return f"openalex_pdf_error: 429 {message[:100]}" return f"openalex_pdf_error: 429 {message[:100]}"