diff --git a/apps/resm/tasks.py b/apps/resm/tasks.py index a3ec563..bbff0e3 100644 --- a/apps/resm/tasks.py +++ b/apps/resm/tasks.py @@ -13,6 +13,7 @@ from datetime import datetime, timedelta import random from .pdf_utils import _is_elsevier_preview_pdf from .d_oaurl import download_from_url_playwright +from uuid import uuid4 import asyncio import sys import os @@ -492,6 +493,24 @@ def touch_alive(def_name: str): """标记该自触发链仍在运行。""" 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): return cache.get(def_name + ":alive") is not None @@ -506,10 +525,13 @@ def ensure_fetch_running(): return f"ensure_fetch_running started: {started}" @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 if not show_task_run(def_name): return "stoped" + chain_id = claim_chain(def_name, chain_id) + if chain_id is None: + return "duplicate chain, exit" touch_alive(def_name) # 限流退避中: 不打 API, 慢节奏自重发只为维持 alive, 等 exceed 标记自然过期。 @@ -517,7 +539,7 @@ def get_pdf_from_openalex(number_of_task: int =10): if cache.get("openalex_api_exceed"): current_app.send_task( "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, ) 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", kwargs={ "number_of_task": number_of_task, + "chain_id": chain_id, }, countdown=countdown, ) @@ -637,13 +660,16 @@ def _elsevier_fetch_pdf(req, paper): @shared_task(base=CustomTask) 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(摘要/全文标记), 若有全文则内联 再取一次 PDF; 并补抓历史上已有全文标记但缺 PDF 的论文。原 get_pdf_from_elsevier 已并入。 number_of_task: 阶段1(摘要+内联 PDF)每轮上限; pdf_number_of_task: 阶段2(存量补 PDF)每轮上限。""" def_name = get_abstract_from_elsevier.name if not show_task_run(def_name): return "stoped" + chain_id = claim_chain(def_name, chain_id) + if chain_id is None: + return "duplicate chain, exit" touch_alive(def_name) # 待抓摘要(并顺带取 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, "exclude_failed": exclude_failed, "pdf_number_of_task": pdf_number_of_task, + "chain_id": chain_id, }, countdown=countdown, ) @@ -870,10 +897,12 @@ def save_pdf_from_openalex(paper:Paper): message = res.json().get("message", "") except ValueError: message = res.text - if "Insufficient credits" in message: + # 文案历史上出现过 "Insufficient credits" 和 "Insufficient budget"(每日额度, + # UTC 午夜重置), 放宽到 "Insufficient" 统一匹配, 避免落进 2 分钟短退避分支空转 + if "Insufficient" in message: # 额度耗尽: 退避 1 小时 cache.set("openalex_api_exceed", True, timeout=3600) - return "openalex_pdf_error: Insufficient credits" + return f"openalex_pdf_error: {message[:100]}" # 普通限流(请求过频): 短退避 2 分钟, 避免立刻重试再撞 429 cache.set("openalex_api_exceed", True, timeout=120) return f"openalex_pdf_error: 429 {message[:100]}"