diff --git a/scripts/parsers/pmc_bulk.py b/scripts/parsers/pmc_bulk.py index 3aa7bfa..e27819d 100644 --- a/scripts/parsers/pmc_bulk.py +++ b/scripts/parsers/pmc_bulk.py @@ -51,6 +51,7 @@ class PMCBulkParser(BaseParser): self, limit: int = 1000, resume_token: str | None = None, + start_after: str | None = None, progress_cb: ProgressCallback | None = None, **_ignored: Any, ) -> Iterator[dict[str, Any]]: @@ -59,6 +60,10 @@ class PMCBulkParser(BaseParser): Args: limit: сколько статей отдать resume_token: продолжить листинг с сохранённой страницы бакета + start_after: ключ, после которого начинать листинг. Нужен, когда + токена нет (первый прогон после ручных заливок): без него + листинг пошёл бы с начала бакета и часами перемалывал уже + залитые статьи как дубли progress_cb: см. base.ProgressCallback """ token = resume_token @@ -66,7 +71,10 @@ class PMCBulkParser(BaseParser): while given < limit: try: - ids, token = self._list_page(token, want=min(200, limit - given)) + ids, token = self._list_page( + token, want=min(200, limit - given), + start_after=None if token else start_after, + ) except Exception as e: logger.error("PMC bulk: листинг не удался: %s", e) return @@ -91,12 +99,16 @@ class PMCBulkParser(BaseParser): if token is None: return - def _list_page(self, token: str | None, want: int) -> tuple[list[str], str | None]: + def _list_page( + self, token: str | None, want: int, start_after: str | None = None + ) -> tuple[list[str], str | None]: """Одна страница листинга бакета: id статей и токен следующей страницы.""" url = f"{BUCKET}/?list-type=2&max-keys=1000" if token: # В токене бывают + и / — без экранирования S3 отвечает 400 url += f"&continuation-token={quote(token, safe='')}" + elif start_after: + url += f"&start-after={quote(start_after, safe='')}" resp = self.client.get(url) resp.raise_for_status() ids = KEY_RE.findall(resp.text)[:want] diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index 95ab943..094e1c0 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -532,6 +532,24 @@ def _stage_work( logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}") +def _last_pmc_key() -> str | None: + """Ключ бакета после последней залитой статьи PMC — старт листинга. + + Бакет отдаётся в лексикографическом порядке ключей, и ext_id вида + `pmc:PMC10000000` соответствует ключу `PMC10000000.1`. Берём максимальный + залитый и продолжаем после него. + """ + from sqlalchemy import text as sql_text + + with db_session() as session: + row = session.execute(sql_text( + "SELECT max(ext_id) FROM documents WHERE ext_id LIKE 'pmc:PMC%'" + )).first() + if not row or not row[0]: + return None + return row[0].split(":", 1)[1] + ".1" + + def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, Any]]: """Инстанс парсера и совместимые с его fetch() аргументы по типу источника.""" limit = cfg["limit"] @@ -571,7 +589,14 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A } elif source_type == "pmc_bulk": from pmc_bulk import PMCBulkParser as P - kwargs = {"limit": limit, "resume_token": cfg.get("resume_token")} + kwargs = { + "limit": limit, + "resume_token": cfg.get("resume_token"), + # Токена ещё нет (первый прогон после ручных заливок) — начинаем + # после последней уже залитой статьи, иначе листинг часами + # перемалывает существующие как дубли + "start_after": None if cfg.get("resume_token") else _last_pmc_key(), + } else: raise ValueError(f"неизвестный тип источника: {source_type}")