#!/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, поэтому прерванная заливка продолжается с той же страницы бакета. Проверено на проде: 1831 отпечаток на статью — реальная глубина, а не аннотация. Темп ~2 статьи/с в один поток (0.3 ст/с из первого замера были получены, когда база одновременно считала эмбеддинги). Скачивание не узкое место (12 статей/с в 12 потоков), упирается в последовательную обработку документа — для миллионов её нужно распараллелить по ядрам воркера. """ import argparse import contextlib 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 write_batch(conn, connect, docs: list[dict], fp_limit: int) -> tuple: """Записать пачку статей, пережив обрыв соединения с базой. Прошлые прогоны умирали целиком от `server closed the connection`, теряя и позицию в бакете. Здесь обрыв стоит переподключения и повтора. Returns: (соединение, добавлено, пропущено, отпечатков) — соединение может быть новым, если пришлось переподключаться """ import psycopg2 from app.algorithms.winnowing import sample_evenly, winnow_ordered for attempt in range(1, 6): try: added = skipped = fp_count = 0 with conn.cursor() as cur: # ON CONFLICT — повторный прогон не плодит дубли 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""", [(d["ext_id"], "pmc", d["title"], d["doi"], d["year"], "en", d["journal"], d["text"][:2000]) for d in docs], ) 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() 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 # Равномерно по тексту, а не первые N: иначе хвост статьи # остаётся без отпечатков hashes = sample_evenly(winnow_ordered(d["text"]), fp_limit) for pos, h in enumerate(hashes): buf.write(f"{doc_id}\t{h}\t{pos}\n") fp_count += len(hashes) added += 1 buf.seek(0) if added: cur.copy_from(buf, "fingerprints", columns=("doc_id", "hash_value", "position")) conn.commit() return conn, added, skipped, fp_count except psycopg2.OperationalError as e: wait = min(60, 5 * attempt) print(f" база недоступна ({str(e).strip()[:60]}), повтор через {wait}с " f"[{attempt}/5]", flush=True) time.sleep(wait) with contextlib.suppress(Exception): conn.close() with contextlib.suppress(Exception): conn = connect() raise RuntimeError("не удалось записать пачку после 5 попыток") 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.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)"}) def connect(): return psycopg2.connect( host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT, dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER, password=settings.POSTGRES_PASSWORD, ) conn = connect() added = skipped = failed = fp_total = done = 0 t0 = time.time() while done < args.limit: want = min(args.batch, args.limit - done) # Сеть у бакета иногда отваливается (прошлый прогон умер на таймауте # SSL-рукопожатия) — это стоит паузы, а не всей заливки try: ids, token = list_articles(client, token, want) except Exception as e: print(f" листинг не удался ({type(e).__name__}), пауза 30с", flush=True) time.sleep(30) continue 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 conn, n_added, n_skipped, n_fp = write_batch(conn, connect, docs, args.fp_per_doc) added += n_added skipped += n_skipped fp_total += n_fp 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} статей/с → " f"миллион за ~{1e6 / max(added / el, 0.01) / 3600:.0f} ч") if __name__ == "__main__": main()