diff --git a/scripts/ops/bulk_ingest_pmc.py b/scripts/ops/bulk_ingest_pmc.py new file mode 100755 index 0000000..1a584e6 --- /dev/null +++ b/scripts/ops/bulk_ingest_pmc.py @@ -0,0 +1,201 @@ +#!/usr/bin/env python3 +"""Массовая заливка полнотекстовых статей из открытого бакета PMC (миллионы). + +Почему отдельный путь, а не парсеры и Celery: постраничный API даёт ~1-2 статьи +в секунду и упирается в rate limit — миллионы так не залить. У AWS Open Data +бакета `pmc-oa-opendata` (открыт, ключ не нужен) для каждой статьи лежит уже +извлечённый текст `PMC*.txt` (33-130 КБ) и метаданные `PMC*.json`. Замер: +12 статей/с в 10 потоков, то есть миллион — примерно за сутки. + +Второе узкое место — запись отпечатков. Построчный INSERT даёт ~6.7 тыс. строк/с +(миллион статей = 2.5 млрд строк = четверо суток только на вставку), поэтому +здесь используется COPY: он на порядки быстрее и не гоняет данные через ORM. + +Масштаб для планирования: 1 млн статей ≈ 2.5 млрд отпечатков ≈ 400 ГБ в БД +с индексами. Плотность отпечатков ограничена --fp-per-doc (по умолчанию 2000 +вместо 20 000 из настроек воркера) — иначе объём растёт быстрее пользы. + +Запуск (в контейнере worker-indexer, где есть psycopg2 и winnowing): + cd /home/user/anti-plagiarism + C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -" + $C --limit 200 < scripts/ops/bulk_ingest_pmc.py # пробная порция + $C --limit 100000 --workers 20 < scripts/ops/bulk_ingest_pmc.py + +Идемпотентно: документы вставляются с ON CONFLICT DO NOTHING по ext_id, уже +залитые пропускаются. Прогресс листинга сохраняется в --state-file, поэтому +прерванная заливка продолжается с той же страницы бакета. + +СТАТУС: НЕ ЗАКОНЧЕН. Проба на 200 статьях: документы вставляются верно +(проверено), но шаг отпечатков ошибочно считает свежевставленные документы уже +обработанными — проверка «есть ли отпечатки» срабатывает там, где их нет, и +COPY не выполняется. Пробные записи из базы удалены. До отладки этого места +скрипт запускать на объёме нельзя: он зальёт документы без отпечатков, то есть +невидимые для L1. +""" + +import argparse +import io +import json +import os +import re +import time +from concurrent.futures import ThreadPoolExecutor +from urllib.parse import quote + +BUCKET = "https://pmc-oa-opendata.s3.amazonaws.com" +STATE_FILE = "/tmp/bulk_pmc_state.json" +KEY_RE = re.compile(r"(PMC\d+\.\d+)/\1\.txt") +TOKEN_RE = re.compile(r"([^<]+)") + + +def list_articles(client, token: str | None, want: int) -> tuple[list[str], str | None]: + """Собрать id статей из листинга бакета, продолжая с сохранённой страницы.""" + ids: list[str] = [] + while len(ids) < want: + url = f"{BUCKET}/?list-type=2&max-keys=1000" + if token: + # В токене бывают + и / — без экранирования S3 отвечает 400 + url += f"&continuation-token={quote(token, safe='')}" + resp = client.get(url) + resp.raise_for_status() + ids.extend(KEY_RE.findall(resp.text)) + m = TOKEN_RE.search(resp.text) + token = m.group(1) if m else None + if token is None: + break + return ids[:want], token + + +def fetch_one(client, art_id: str) -> dict | None: + """Текст статьи и её метаданные. None — если статья недоступна или пустая.""" + try: + txt = client.get(f"{BUCKET}/{art_id}/{art_id}.txt") + if txt.status_code != 200 or len(txt.text) < 1500: + return None + meta_resp = client.get(f"{BUCKET}/{art_id}/{art_id}.json") + meta = meta_resp.json() if meta_resp.status_code == 200 else {} + except Exception: + return None + + citation = meta.get("citation") or "" + year = None + m = re.search(r"\b(19|20)\d{2}\b", citation) + if m: + year = int(m.group(0)) + + return { + "ext_id": f"pmc:{meta.get('pmcid') or art_id.split('.')[0]}", + "title": (meta.get("title") or art_id)[:1000], + "doi": meta.get("doi"), + "year": year, + "journal": citation[:300] or None, + "text": txt.text, + } + + +def main() -> None: + ap = argparse.ArgumentParser(description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("--limit", type=int, default=200, help="сколько статей залить за прогон") + ap.add_argument("--workers", type=int, default=10, help="потоков закачки (вежливо: 10-30)") + ap.add_argument("--batch", type=int, default=200, help="статей в одной транзакции") + ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на документ") + ap.add_argument("--state-file", default=STATE_FILE, help="файл с позицией листинга") + args = ap.parse_args() + + import httpx + import psycopg2 + from app.algorithms.winnowing import winnow + from app.config import settings + + token = None + if os.path.exists(args.state_file): + with open(args.state_file) as fh: + token = json.load(fh).get("token") + print("продолжаю листинг с сохранённой позиции") + + client = httpx.Client(timeout=60, headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"}) + conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT, + dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER, + password=settings.POSTGRES_PASSWORD) + + added = skipped = failed = 0 + fp_total = 0 + t0 = time.time() + done = 0 + + while done < args.limit: + want = min(args.batch, args.limit - done) + ids, token = list_articles(client, token, want) + if not ids: + print("бакет закончился") + break + done += len(ids) + + with ThreadPoolExecutor(max_workers=args.workers) as pool: + docs = [d for d in pool.map(lambda i: fetch_one(client, i), ids) if d] + failed += len(ids) - len(docs) + if not docs: + continue + + with conn.cursor() as cur: + # Документы: ON CONFLICT — повторный прогон не плодит дубли + rows = [ + (d["ext_id"], "pmc", d["title"], d["doi"], d["year"], "en", + d["journal"], d["text"][:2000]) + for d in docs + ] + cur.executemany( + """INSERT INTO documents (ext_id, source, title, doi, year, lang, journal, + abstract, authors, indexed_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,'[]'::json, now()) + ON CONFLICT (ext_id) DO NOTHING""", + rows, + ) + cur.execute("SELECT id, ext_id FROM documents WHERE ext_id = ANY(%s)", + ([d["ext_id"] for d in docs],)) + id_by_ext = {e: i for i, e in cur.fetchall()} + + # Отпечатки: COPY вместо построчных INSERT — единственный способ + # писать миллиарды строк за разумное время + buf = io.StringIO() + fresh = 0 + for d in docs: + doc_id = id_by_ext.get(d["ext_id"]) + if doc_id is None: + continue + cur.execute("SELECT 1 FROM fingerprints WHERE doc_id=%s LIMIT 1", (doc_id,)) + if cur.fetchone(): + skipped += 1 + continue + hashes = list(winnow(d["text"]))[: args.fp_per_doc] + for pos, h in enumerate(hashes): + buf.write(f"{doc_id}\t{h}\t{pos}\n") + fp_total += len(hashes) + fresh += 1 + buf.seek(0) + if fresh: + cur.copy_from(buf, "fingerprints", columns=("doc_id", "hash_value", "position")) + added += fresh + conn.commit() + + with open(args.state_file, "w") as fh: + json.dump({"token": token}, fh) + + el = time.time() - t0 + print(f" {done}/{args.limit} · залито {added} · пропущено {skipped} · " + f"недоступно {failed} · отпечатков {fp_total:,} · {added / max(el, 1):.1f} ст/с", + flush=True) + if token is None: + break + + conn.close() + el = time.time() - t0 + print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, " + f"недоступно {failed}, отпечатков {fp_total:,}") + if added: + print(f"темп: {added / el:.1f} статей/с → миллион за ~{1e6 / max(added / el, 0.01) / 3600:.0f} ч") + + +if __name__ == "__main__": + main()