#!/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, поэтому прерванная заливка продолжается с той же страницы бакета. Проверено на проде 2026-08-31: 100 статей залито, 183 090 отпечатков (1831 на статью — реальная глубина, а не аннотация). Темп — 0.3 статьи/с, то есть миллион в один поток занял бы около 40 суток: скачивание тут не узкое место (12 статей/с в 12 потоков), упирается в последовательную обработку одного документа. Для миллионов обработку нужно распараллелить по ядрам воркера — это следующий шаг, до него скрипт годится для порций в десятки тысяч. """ 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()