#!/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()