From 26de1b9de101377165efe34d46e5ea7b07515b4f Mon Sep 17 00:00:00 2001 From: jze9 Date: Sat, 5 Sep 2026 20:32:53 +0500 Subject: [PATCH] =?UTF-8?q?refactor(ops):=20=D0=BC=D0=B0=D1=81=D1=81=D0=BE?= =?UTF-8?q?=D0=B2=D0=B0=D1=8F=20=D0=B7=D0=B0=D0=BB=D0=B8=D0=B2=D0=BA=D0=B0?= =?UTF-8?q?=20=E2=80=94=20=D0=B2=20=D0=BE=D0=B1=D1=89=D0=B8=D0=B9=20=D0=BA?= =?UTF-8?q?=D0=BE=D0=BD=D0=B2=D0=B5=D0=B9=D0=B5=D1=80=20=D0=B2=D0=BC=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D0=BE=20=D0=BE=D1=82=D0=B4=D0=B5=D0=BB=D1=8C=D0=BD?= =?UTF-8?q?=D1=8B=D1=85=20=D1=81=D0=BA=D1=80=D0=B8=D0=BF=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Заливку Википедии и PMC я сделал отдельными скриптами мимо существующей инфраструктуры: запуск руками через ssh, состояние в файле в /tmp, никакой видимости. Результат предсказуем — за неделю обе умерли молча (обрыв базы на 20 008 статьях из 2 млн и таймаут сети на 85 тыс. из 100 тыс.), прогресс потерялся, а узнали мы об этом через неделю. При том что рядом лежит готовый механизм: parse_sources, прогоны со шкалой, журнал, кнопки, ретраи Celery. Теперь это обычные типы источника — wikipedia_ru и pmc_bulk: - заводятся и запускаются из админки, как OpenAlex или КиберЛенинка; - показывают ту же шкалу, журнал и кнопку остановки; - падение воркера больше не теряет прогресс: позиция продолжения хранится в parse_sources.resume_token (номер статьи в дампе / токен страницы бакета), повторный запуск берёт следующую порцию; - укладываются в бюджет времени таска — заливка идёт порциями, а не одним многосуточным процессом. Чего не хватало конвейеру для миллионов и что добавлено: - парсеры отдают генератор, а не список: 2 млн статей в память не влезают; - app/bulk_writer.py — запись пачками через COPY (21 тыс. строк/с против 6.7 тыс. построчно) с переподключением к базе при обрыве; - эмбеддинги при массовой заливке не диспатчатся: они на порядок медленнее и стали бы узким местом, вектора досчитываются отдельно (reembed_missing.py). scripts/ops/bulk_ingest_*.py удалены — их работу делает конвейер. Проверено на проде: оба парсера отдают документы, прогон через run_parser завершается штатно, позиция продолжения сдвигается (300 → 600), повторный запуск продолжает с неё. Co-Authored-By: Claude Opus 5 --- docs/ARCHITECTURE.md | 22 +- docs/INGESTION.md | 39 ++- scripts/ops/bulk_ingest_pmc.py | 247 ---------------- scripts/ops/bulk_ingest_wikipedia_ru.py | 269 ------------------ scripts/parsers/pmc_bulk.py | 144 ++++++++++ scripts/parsers/wikipedia_ru.py | 179 ++++++++++++ .../versions/007_parse_source_resume.py | 29 ++ services/api/app/api/admin.py | 6 +- services/api/app/models/admin.py | 3 + services/api/app/schemas/admin.py | 8 +- services/frontend/src/pages/admin/Sources.tsx | 6 +- services/worker-indexer/app/bulk_writer.py | 132 +++++++++ services/worker-indexer/app/config.py | 6 + .../worker-indexer/app/models/__init__.py | 1 + services/worker-indexer/app/tasks/index.py | 225 ++++++++++++--- 15 files changed, 732 insertions(+), 584 deletions(-) delete mode 100755 scripts/ops/bulk_ingest_pmc.py delete mode 100755 scripts/ops/bulk_ingest_wikipedia_ru.py create mode 100644 scripts/parsers/pmc_bulk.py create mode 100644 scripts/parsers/wikipedia_ru.py create mode 100644 services/api/alembic/versions/007_parse_source_resume.py create mode 100644 services/worker-indexer/app/bulk_writer.py diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index eeff53e..bdbfac7 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -78,7 +78,9 @@ - **fingerprints** — `doc_id`, `hash_value` (BIGINT, Winnowing), `position` — для L1. - **usage_logs** — `user_id`, `action` — учёт лимитов по тарифу. - **parse_sources** — задания парсеров (админка): тип, query, годы, лимит, статус, - `last_run_id` — ссылка на последний прогон. + `last_run_id` — ссылка на последний прогон, `resume_token` — позиция + продолжения для массовых источников (номер статьи в дампе, токен страницы + бакета). - **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/ added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал событий (JSON). Тайминги раздельные: `started_at` — постановка в очередь, @@ -119,10 +121,20 @@ `app.bibliography.build_bibliography` (сортировка кириллица→латиница, нумерация, формат 7.1/7.0.5) → результат. -**Наполнение корпуса.** админка/CLI → `index.run_parser` (OpenAlex/arXiv/PMC/ -КиберЛенинка, фильтр `is_oa`) → `index.add_document` (дедуп по `ext_id`, -fingerprints, MinHash, эмбеддинг) → `index.enrich_full_text` (скачать OA-PDF → -MinIO → переиндексация). Ход заливки пишется в `parse_runs` (§12). +**Наполнение корпуса.** админка → `index.run_parser` — один конвейер для всех +источников, со шкалой, журналом и кнопкой остановки. Внутри два пути записи: + +- *обычные источники* (OpenAlex/arXiv/PMC/КиберЛенинка, фильтр `is_oa`) — + `index.add_document` на каждый документ: дедуп по `ext_id`, fingerprints, + MinHash, Elasticsearch, эмбеддинг, затем `index.enrich_full_text` (скачать + OA-PDF → MinIO → переиндексация); +- *массовые источники* (`wikipedia_ru`, `pmc_bulk`) — парсер отдаёт генератор, + запись идёт пачками через `COPY` (`app/bulk_writer.py`), позиция продолжения + хранится в `parse_sources.resume_token`. Эмбеддинги там не считаются: они + медленнее заливки на порядок и стали бы её узким местом, вектора + досчитываются отдельно (§12). + +Ход заливки в обоих случаях пишется в `parse_runs` (§12). Второй путь наполнения — ручная загрузка файлов админом: api сохраняет их в MinIO (`corpus-upload/`) → `index.ingest_upload` (извлечь текст → `add_document`, `source=manual_upload`). Это не проверка на плагиат: файл сразу становится diff --git a/docs/INGESTION.md b/docs/INGESTION.md index 0f58d86..ce01ede 100644 --- a/docs/INGESTION.md +++ b/docs/INGESTION.md @@ -62,7 +62,7 @@ | Источник | Объём | Годен для массовой заливки | |----------|-------|----------------------------| -| **Википедия ru** | дамп 5.6 ГБ, ~2 млн статей | **да** — `bulk_ingest_wikipedia_ru.py`, потоком без блокировок, ~919 отпечатков на статью. Требует осмысленный User-Agent, иначе 403 | +| **Википедия ru** | дамп 5.6 ГБ, ~2 млн статей | **да** — тип источника `wikipedia_ru`, ~919 отпечатков на статью. Дамп качается заранее (Wikimedia отдаёт 403 без осмысленного User-Agent и рвёт долгие потоковые соединения) | | КиберЛенинка | ~3 млн статей | нет — блокирует выкачку PDF после ~130 запросов | | OpenAlex `language:ru` + OA | заявлено 395 410 | практически нет — по прямым `pdf_url` скачалось 2 из 10 (остальное 403 издателей), а метка языка ненадёжна: в выдаче попадаются англоязычные журналы | | eLIBRARY.RU (РИНЦ) | ~40 млн | только по договору, публичной выгрузки нет | @@ -72,19 +72,36 @@ Википедия формально не научный источник, но студенты копируют из неё чаще всего, а по объёму связного русского текста ей нет альтернативы среди доступного. -### Массовая заливка: bulk вместо API +### Массовые источники — тот же конвейер, что и обычные Постраничные API дают 1-2 статьи в секунду и упираются в rate limit — миллионы -так не залить. Для объёма есть `scripts/ops/bulk_ingest_pmc.py`: открытый бакет -`pmc-oa-opendata` (ключ не нужен) отдаёт у каждой статьи готовый извлечённый -текст. Проверено 31.08: 100 статей, 183 090 отпечатков — 1831 на статью, то есть -настоящая глубина, а не аннотация. +так не залить. Для объёма есть два типа источника, которые читают дамп или +бакет потоком: + +| Тип | Откуда | Особенность | +|-----|--------|-------------| +| `wikipedia_ru` | локальный дамп `scripts/parsers/ruwiki.xml.bz2` | позиция продолжения — номер статьи в дампе | +| `pmc_bulk` | бакет `pmc-oa-opendata` (открыт, ключ не нужен) | позиция — токен страницы бакета; у каждой статьи готовый извлечённый текст | + +Заводятся и запускаются они как любой другой источник — через админку, и точно +так же показывают шкалу, журнал и кнопку остановки. Отличия внутри: + +- парсер отдаёт **генератор**, а не список: миллионы статей в память не влезут; +- запись идёт пачками через `COPY` (`worker-indexer/app/bulk_writer.py`) — + построчная вставка даёт 6.7 тыс. строк/с против 21 тыс. у COPY, а миллион + статей это ~2 млрд отпечатков; +- **эмбеддинги при заливке не считаются**: они медленнее заливки и сделали бы её + узким местом. L1 и L2 работают сразу, векторы досчитываются потом + (`scripts/ops/reembed_missing.py`); +- позиция продолжения хранится в `parse_sources.resume_token`, поэтому источник + запускается повторно до исчерпания — каждый прогон берёт следующую порцию и + укладывается в бюджет времени таска. + +Планировать объём: 1 млн статей ≈ 2 млрд отпечатков ≈ 200 ГБ в базе с индексами. +При таком росте индексы перестают помещаться в память сервера БД — проверено на +практике: поиск L1 деградировал с 31 мс до 484 мс, пока не увеличили RAM и +`shared_buffers` (см. DR-HA.md). -Ограничение по темпу: 0.3 статьи/с, миллион в один поток — около 40 суток. -Скачивание тут не узкое место (12 статей/с в 12 потоков), упирается в -последовательную обработку документа; для миллионов её надо распараллелить по -ядрам воркера. Планировать объём: 1 млн статей ≈ 2.5 млрд отпечатков ≈ 400 ГБ -в базе с индексами (сейчас 113 млн строк занимают 12 ГБ). ## Что подготовлено - **Парсер CyberLeninka починен** (`scripts/parsers/cyberleninka.py`): раньше слал GET diff --git a/scripts/ops/bulk_ingest_pmc.py b/scripts/ops/bulk_ingest_pmc.py deleted file mode 100755 index 042b229..0000000 --- a/scripts/ops/bulk_ingest_pmc.py +++ /dev/null @@ -1,247 +0,0 @@ -#!/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() diff --git a/scripts/ops/bulk_ingest_wikipedia_ru.py b/scripts/ops/bulk_ingest_wikipedia_ru.py deleted file mode 100755 index 7625703..0000000 --- a/scripts/ops/bulk_ingest_wikipedia_ru.py +++ /dev/null @@ -1,269 +0,0 @@ -#!/usr/bin/env python3 -"""Массовая заливка русской Википедии в корпус сравнения. - -Зачем именно она: студенты копируют из Википедии чаще, чем из научных статей, а -в корпусе её нет вовсе. При этом она единственный русскоязычный источник такого -объёма, который отдаётся без блокировок — КиберЛенинка режет выкачку после ~130 -запросов, eLIBRARY требует договора. Дамп `ruwiki-latest-pages-articles.xml.bz2` -(5.6 ГБ, ~2 млн статей) скачивается свободно. - -Дамп не сохраняется на диск: читается потоком и распаковывается на лету — -на app-хосте всего 22 ГБ свободно. - -Что отбрасывается: перенаправления, служебные пространства имён (оставляем -только статьи) и короткие тексты — для детекции нужен объём, а не заготовки. - -Запуск (в контейнере worker-indexer): - cd /home/user/anti-plagiarism - C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -" - $C --limit 100 < scripts/ops/bulk_ingest_wikipedia_ru.py # проба - $C --limit 500000 < scripts/ops/bulk_ingest_wikipedia_ru.py # объём - -Идемпотентно: ext_id вида `wikipedia_ru:`, повторная заливка пропускает -уже существующие статьи (ON CONFLICT DO NOTHING). -""" - -import argparse -import bz2 -import contextlib -import io -import json -import os -import re -import time - -DUMP_URL = "https://dumps.wikimedia.org/ruwiki/latest/ruwiki-latest-pages-articles.xml.bz2" -PAGE_RE = re.compile(r"(.*?)", re.DOTALL) -TITLE_RE = re.compile(r"(.*?)", re.DOTALL) -ID_RE = re.compile(r"(\d+)") -NS_RE = re.compile(r"(\d+)") -TEXT_RE = re.compile(r']*>(.*?)', re.DOTALL) -REDIRECT_RE = re.compile(r"]*>.*?", re.DOTALL), " "), - (re.compile(r"<[^>]+>"), " "), # html-теги - (re.compile(r"^[*#:;|!].*$", re.MULTILINE), " "), # списки и таблицы - (re.compile(r"^=+.*?=+$", re.MULTILINE), " "), # == заголовки разделов == - (re.compile(r"'{2,}"), ""), # ''курсив'' - (re.compile(r"&[a-z]+;"), " "), - (re.compile(r"[ \t]+"), " "), - (re.compile(r"\n{2,}"), "\n"), -] - - -def clean_wikitext(raw: str) -> str: - text = raw - for _ in range(3): # шаблоны бывают вложенными - text = CLEAN_RULES[0][0].sub(CLEAN_RULES[0][1], text) - for rx, repl in CLEAN_RULES[1:]: - text = rx.sub(repl, text) - return text.strip() - - -def _pages_from_chunks(chunks, min_chars: int): - """Разобрать поток распакованного XML на статьи основного пространства имён.""" - buf = "" - for raw in chunks: - buf += raw.decode("utf-8", errors="replace") - - while True: - m = PAGE_RE.search(buf) - if not m: - break - page, buf = m.group(1), buf[m.end():] - - if REDIRECT_RE.search(page): - continue - ns = NS_RE.search(page) - if not ns or ns.group(1) != "0": # только статьи - continue - tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page) - if not (tm and im and xm): - continue - text = clean_wikitext(xm.group(1)) - if len(text) < min_chars: - continue - yield im.group(1), tm.group(1), text - - if len(buf) > 20 * 1024 * 1024: # страховка от разбухания - buf = buf[-1024 * 1024:] - - -def iter_pages(min_chars: int, dump_file: str | None): - """Статьи из локального дампа либо, если файла нет, прямо из сети. - - Локальный файл предпочтителен: Wikimedia рвёт долгие потоковые соединения - (проверено — обрыв на 32 МБ из 5.9 ГБ), а скачать файл можно с докачкой - (`curl -C -`) и потом читать сколько угодно. - """ - decomp = bz2.BZ2Decompressor() - - if dump_file: - def chunks(): - with open(dump_file, "rb") as fh: - while True: - part = fh.read(4 * 1024 * 1024) - if not part: - return - try: - yield decomp.decompress(part) - except EOFError: - return - yield from _pages_from_chunks(chunks(), min_chars) - return - - import httpx - - # Wikimedia отдаёт 403 без осмысленного User-Agent — по их правилам он должен - # называть приложение и давать контакт - headers = {"User-Agent": "AcademicHelper/1.0 (https://academic.jze9.ru; noreply@jze9.ru)"} - - def net_chunks(): - with httpx.stream("GET", DUMP_URL, timeout=120, follow_redirects=True, - headers=headers) as resp: - resp.raise_for_status() - for chunk in resp.iter_bytes(4 * 1024 * 1024): - try: - yield decomp.decompress(chunk) - except EOFError: - return - - yield from _pages_from_chunks(net_chunks(), min_chars) - - -def main() -> None: - ap = argparse.ArgumentParser(description=__doc__, - formatter_class=argparse.RawDescriptionHelpFormatter) - ap.add_argument("--limit", type=int, default=100, help="сколько статей залить") - ap.add_argument("--batch", type=int, default=150, help="статей в транзакции") - ap.add_argument("--min-chars", type=int, default=2000, help="минимальная длина текста") - ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на статью") - ap.add_argument("--dump-file", default=None, - help="путь к заранее скачанному дампу (надёжнее, чем читать из сети)") - ap.add_argument("--state-file", default="/tmp/bulk_wiki_state.json", - help="файл с числом уже обработанных статей — для продолжения после сбоя") - args = ap.parse_args() - - import psycopg2 - from app.algorithms.winnowing import sample_evenly, winnow_ordered - from app.config import settings - - 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 = 0 - fp_total = 0 - seen = 0 # сколько статей прошло через парсер (позиция в дампе) - t0 = time.time() - batch: list[tuple[str, str, str]] = [] - - # Позиция в дампе: прошлый прогон обрывался на 20 008 статье и начинал всё - # сначала — дедуп по ext_id спасал от дублей, но время тратилось впустую - skip_until = 0 - if os.path.exists(args.state_file): - with open(args.state_file) as fh: - skip_until = json.load(fh).get("processed", 0) - if skip_until: - print(f"продолжаю с {skip_until}-й статьи дампа") - - def _flush_once(batch): - """Одна попытка записать пачку; счётчики трогает только при успехе.""" - nonlocal added, skipped, fp_total - with conn.cursor() as cur: - cur.executemany( - """INSERT INTO documents (ext_id, source, title, lang, abstract, url, - authors, indexed_at) - VALUES (%s, 'wikipedia_ru', %s, 'ru', %s, %s, '[]'::json, now()) - ON CONFLICT (ext_id) DO NOTHING""", - [(f"wikipedia_ru:{pid}", title[:1000], text[:2000], - f"https://ru.wikipedia.org/?curid={pid}") for pid, title, text in batch], - ) - cur.execute("SELECT id, ext_id FROM documents WHERE ext_id = ANY(%s)", - ([f"wikipedia_ru:{pid}" for pid, _, _ in batch],)) - id_by_ext = {e: i for i, e in cur.fetchall()} - - buf = io.StringIO() - fresh = 0 - for pid, _, text in batch: - doc_id = id_by_ext.get(f"wikipedia_ru:{pid}") - 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(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() - - def flush(batch): - """Записать пачку, пережив обрыв соединения с базой. - - Прошлый прогон умер целиком на `server closed the connection` — потеряв - и позицию в дампе. Здесь обрыв стоит переподключения и повтора. - """ - nonlocal conn - if not batch: - return - for attempt in range(1, 6): - try: - _flush_once(batch) - return - 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() - try: - conn = connect() - except Exception: - continue - raise RuntimeError("не удалось записать пачку после 5 попыток") - - for pid, title, text in iter_pages(args.min_chars, args.dump_file): - seen += 1 - if seen <= skip_until: # быстро проматываем уже обработанное - continue - batch.append((pid, title, text)) - if len(batch) >= args.batch: - flush(batch) - batch = [] - with open(args.state_file, "w") as fh: - json.dump({"processed": seen}, fh) - el = time.time() - t0 - print(f" залито {added} · пропущено {skipped} · отпечатков {fp_total:,} · " - f"{added / max(el, 1):.1f} ст/с", flush=True) - if added >= args.limit: - break - flush(batch) - conn.close() - - el = time.time() - t0 - print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, " - f"отпечатков {fp_total:,}") - - -if __name__ == "__main__": - main() diff --git a/scripts/parsers/pmc_bulk.py b/scripts/parsers/pmc_bulk.py new file mode 100644 index 0000000..3aa7bfa --- /dev/null +++ b/scripts/parsers/pmc_bulk.py @@ -0,0 +1,144 @@ +"""Парсер PubMed Central из открытого бакета AWS — массовый путь. + +Отличие от `pmc.py`: тот ходит в API E-utilities и годится для точечной +докачки по теме (1-2 статьи в секунду, есть rate limit). Этот берёт статьи из +бакета `pmc-oa-opendata`, где у каждой лежит **уже извлечённый текст**, и +качает их пачками в несколько потоков — путь к сотням тысяч и миллионам. + +Ключ для доступа не нужен, бакет открыт. Замер: 12 статей/с в 10 потоков. + +Как и парсер Википедии, отдаёт статьи генератором: складывать миллионы +документов в список нельзя (см. bulk_writer.py). +""" + +import logging +import re +from collections.abc import Iterator +from concurrent.futures import ThreadPoolExecutor +from typing import Any +from urllib.parse import quote + +import httpx +from base import BaseParser, ProgressCallback + +logger = logging.getLogger(__name__) + +BUCKET = "https://pmc-oa-opendata.s3.amazonaws.com" +KEY_RE = re.compile(r"(PMC\d+\.\d+)/\1\.txt") +TOKEN_RE = re.compile(r"([^<]+)") +YEAR_RE = re.compile(r"\b(19|20)\d{2}\b") +MIN_CHARS = 1500 # меньше — обрывок, а не статья + + +class PMCBulkParser(BaseParser): + """PMC Open Access из бакета AWS. Отдаёт статьи потоком.""" + + source_name = "pmc" # тот же источник в корпусе, что и у API-парсера + bulk = True + + def __init__(self, workers: int = 12) -> None: + super().__init__() + self.workers = workers + self.client = httpx.Client( + timeout=60, + headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"}, + ) + # Позиция листинга: бакет отдаётся страницами, и продолжать прогон + # нужно с той же страницы, а не с начала + self.last_token: str | None = None + + def fetch( # type: ignore[override] + self, + limit: int = 1000, + resume_token: str | None = None, + progress_cb: ProgressCallback | None = None, + **_ignored: Any, + ) -> Iterator[dict[str, Any]]: + """Статьи из бакета: листинг страницами, скачивание в потоках. + + Args: + limit: сколько статей отдать + resume_token: продолжить листинг с сохранённой страницы бакета + progress_cb: см. base.ProgressCallback + """ + token = resume_token + given = 0 + + while given < limit: + try: + ids, token = self._list_page(token, want=min(200, limit - given)) + except Exception as e: + logger.error("PMC bulk: листинг не удался: %s", e) + return + if not ids: + logger.info("PMC bulk: бакет закончился") + return + + with ThreadPoolExecutor(max_workers=self.workers) as pool: + for raw in pool.map(self._fetch_article, ids): + if raw is None: + continue + given += 1 + raw["resume_token"] = token + yield raw + if given >= limit: + break + + self.last_token = token + if progress_cb and not progress_cb(given): + logger.info("PMC bulk: выборка остановлена по запросу (%d)", given) + return + if token is None: + return + + def _list_page(self, token: str | None, want: int) -> 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='')}" + resp = self.client.get(url) + resp.raise_for_status() + ids = KEY_RE.findall(resp.text)[:want] + m = TOKEN_RE.search(resp.text) + return ids, (m.group(1) if m else None) + + def _fetch_article(self, art_id: str) -> dict[str, Any] | None: + """Текст статьи и метаданные; None — статья недоступна или пустая.""" + try: + txt = self.client.get(f"{BUCKET}/{art_id}/{art_id}.txt") + if txt.status_code != 200 or len(txt.text) < MIN_CHARS: + return None + meta_resp = self.client.get(f"{BUCKET}/{art_id}/{art_id}.json") + meta = meta_resp.json() if meta_resp.status_code == 200 else {} + except Exception: + return None + + return {"art_id": art_id, "meta": meta, "text": txt.text} + + def transform(self, raw: dict[str, Any]) -> dict[str, Any]: + """Привести статью к унифицированному формату документов корпуса.""" + meta = raw.get("meta") or {} + art_id = raw.get("art_id", "") + pmcid = meta.get("pmcid") or art_id.split(".")[0] + if not pmcid: + return {} + + citation = meta.get("citation") or "" + year_match = YEAR_RE.search(citation) + text = raw.get("text", "") + + return { + "source": self.source_name, + "ext_id": f"pmc:{pmcid}", + "title": (meta.get("title") or pmcid)[:1000], + "authors": [], + "doi": meta.get("doi"), + "year": int(year_match.group(0)) if year_match else None, + "lang": "en", + "journal": citation[:300] or None, + "url": f"https://www.ncbi.nlm.nih.gov/pmc/articles/{pmcid}/", + "abstract": text[:2000], + "text": text, + "resume_token": raw.get("resume_token"), + } diff --git a/scripts/parsers/wikipedia_ru.py b/scripts/parsers/wikipedia_ru.py new file mode 100644 index 0000000..5bcea3a --- /dev/null +++ b/scripts/parsers/wikipedia_ru.py @@ -0,0 +1,179 @@ +"""Парсер русской Википедии из дампа Wikimedia. + +Зачем в корпусе: студенты копируют из Википедии чаще, чем из научных статей, а +по объёму связного русского текста ей нет альтернативы среди доступного — +КиберЛенинка блокирует выкачку, eLIBRARY требует договора. + +Отличие от остальных парсеров: `fetch` возвращает **генератор**, а не список. +Дамп — 5.6 ГБ и около 2 млн статей, держать их в памяти нельзя, поэтому +`run_parser` читает результат лениво и пишет пачками (см. bulk_writer.py). + +Дамп качается заранее и кладётся туда, где его видит воркер: + curl -L -C - --retry 100 -A "AcademicHelper/1.0 (noreply@jze9.ru)" \\ + -o scripts/parsers/ruwiki.xml.bz2 \\ + https://dumps.wikimedia.org/ruwiki/latest/ruwiki-latest-pages-articles.xml.bz2 + +Читать дамп прямо из сети не выйдет: Wikimedia обрывает долгие соединения +(проверено — обрыв на 32 МБ из 5.9 ГБ), да и продолжить с места разрыва +распакованный поток bz2 нельзя. +""" + +import bz2 +import logging +import os +import re +from collections.abc import Iterator +from typing import Any + +from base import BaseParser, ProgressCallback + +logger = logging.getLogger(__name__) + +DEFAULT_DUMP = "/parsers/ruwiki.xml.bz2" + +PAGE_RE = re.compile(r"(.*?)", re.DOTALL) +TITLE_RE = re.compile(r"(.*?)", re.DOTALL) +ID_RE = re.compile(r"(\d+)") +NS_RE = re.compile(r"(\d+)") +TEXT_RE = re.compile(r"]*>(.*?)", re.DOTALL) +REDIRECT_RE = re.compile(r"]*>.*?", re.DOTALL), " "), + (re.compile(r"<[^>]+>"), " "), + (re.compile(r"^[*#:;|!].*$", re.MULTILINE), " "), # списки и таблицы + (re.compile(r"^=+.*?=+$", re.MULTILINE), " "), # == заголовки разделов == + (re.compile(r"'{2,}"), ""), + (re.compile(r"&[a-z]+;"), " "), + (re.compile(r"[ \t]+"), " "), + (re.compile(r"\n{2,}"), "\n"), +] + + +def clean_wikitext(raw: str) -> str: + """Убрать вики-разметку, оставив читаемый текст статьи.""" + text = raw + for _ in range(3): # шаблоны бывают вложенными + text = TEMPLATE_RE.sub(" ", text) + for rx, repl in CLEAN_RULES: + text = rx.sub(repl, text) + return text.strip() + + +class WikipediaRuParser(BaseParser): + """Русская Википедия из локального дампа. Отдаёт статьи потоком.""" + + source_name = "wikipedia_ru" + # Признак для run_parser: результат читается лениво и пишется пачками, + # а не собирается в список и не идёт через add_document по одному + bulk = True + + def fetch( # type: ignore[override] + self, + limit: int = 1000, + min_chars: int = 2000, + dump_path: str | None = None, + skip: int = 0, + progress_cb: ProgressCallback | None = None, + **_ignored: Any, + ) -> Iterator[dict[str, Any]]: + """Статьи основного пространства имён из дампа. + + Args: + limit: сколько статей отдать + min_chars: минимальная длина текста — заготовки в корпусе бесполезны + dump_path: путь к дампу (по умолчанию /parsers/ruwiki.xml.bz2) + skip: сколько статей пропустить — так прогон продолжается с места + обрыва, не перечитывая уже залитое + progress_cb: см. base.ProgressCallback + """ + path = dump_path or os.environ.get("WIKIPEDIA_DUMP_PATH") or DEFAULT_DUMP + if not os.path.exists(path): + raise FileNotFoundError( + f"дамп Википедии не найден: {path} — скачайте его (см. модуль) " + f"или укажите WIKIPEDIA_DUMP_PATH" + ) + + decomp = bz2.BZ2Decompressor() + buf = "" + seen = 0 # статей прошло через парсер — позиция в дампе + given = 0 # статей отдано наружу + + with open(path, "rb") as fh: + while given < limit: + part = fh.read(4 * 1024 * 1024) + if not part: + break + try: + raw = decomp.decompress(part) + except EOFError: + break + if not raw: + continue + buf += raw.decode("utf-8", errors="replace") + + while given < limit: + m = PAGE_RE.search(buf) + if not m: + break + page, buf = m.group(1), buf[m.end():] + + doc = self._page_to_doc(page, min_chars) + if doc is None: + continue + + seen += 1 + if seen <= skip: # быстро проматываем уже обработанное + continue + + given += 1 + doc["dump_position"] = seen + yield doc + + if progress_cb and not progress_cb(given): + logger.info("Википедия: выборка остановлена по запросу (%d)", given) + return + + if len(buf) > 20 * 1024 * 1024: # страховка от разбухания + buf = buf[-1024 * 1024:] + + def _page_to_doc(self, page: str, min_chars: int) -> dict[str, Any] | None: + """Разобрать в документ; None — страница нам не подходит.""" + if REDIRECT_RE.search(page): + return None + ns = NS_RE.search(page) + if not ns or ns.group(1) != "0": # только статьи + return None + + tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page) + if not (tm and im and xm): + return None + + text = clean_wikitext(xm.group(1)) + if len(text) < min_chars: + return None + + return {"pageid": im.group(1), "title": tm.group(1), "text": text} + + def transform(self, raw: dict[str, Any]) -> dict[str, Any]: + """Привести статью к унифицированному формату документов корпуса.""" + pid = raw.get("pageid") + if not pid: + return {} + return { + "source": self.source_name, + "ext_id": f"wikipedia_ru:{pid}", + "title": raw.get("title", ""), + "authors": [], + "year": None, + "lang": "ru", + "url": f"https://ru.wikipedia.org/?curid={pid}", + "abstract": (raw.get("text") or "")[:2000], + "text": raw.get("text", ""), + "dump_position": raw.get("dump_position"), + } diff --git a/services/api/alembic/versions/007_parse_source_resume.py b/services/api/alembic/versions/007_parse_source_resume.py new file mode 100644 index 0000000..e3f109b --- /dev/null +++ b/services/api/alembic/versions/007_parse_source_resume.py @@ -0,0 +1,29 @@ +"""Позиция продолжения массовой заливки в источнике. + +Revision ID: 007 +Revises: 006 +Create Date: 2026-09-05 + +Массовые источники (дамп Википедии, бакет PMC) заливаются частями: за один +прогон берётся столько статей, сколько влезает в бюджет времени таска. Чтобы +следующий прогон продолжал с того же места, а не перечитывал залитое заново, +позиция сохраняется прямо в источнике: для Википедии это номер статьи в дампе, +для PMC — токен страницы бакета. +""" + +import sqlalchemy as sa + +from alembic import op + +revision = "007" +down_revision = "006" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.add_column("parse_sources", sa.Column("resume_token", sa.Text(), nullable=True)) + + +def downgrade() -> None: + op.drop_column("parse_sources", "resume_token") diff --git a/services/api/app/api/admin.py b/services/api/app/api/admin.py index d3273c9..52b00b5 100644 --- a/services/api/app/api/admin.py +++ b/services/api/app/api/admin.py @@ -50,8 +50,10 @@ logger = logging.getLogger(__name__) router = APIRouter(prefix="/admin", tags=["admin"], dependencies=[Depends(get_admin_user)]) -# Типы источников, для которых в worker-indexer есть парсер (index.run_parser) -SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc") +# Типы источников, для которых в worker-indexer есть парсер (index.run_parser). +# wikipedia_ru и pmc_bulk — массовые: читают дамп/бакет потоком и пишут пачками, +# продолжая с сохранённой позиции (parse_sources.resume_token) +SOURCE_TYPES = ("openalex", "cyberleninka", "arxiv", "pmc", "wikipedia_ru", "pmc_bulk") # Прогон в этих статусах ещё может двигаться — повторно запускать источник нельзя RUN_ACTIVE_STATUSES = ("queued", "running") diff --git a/services/api/app/models/admin.py b/services/api/app/models/admin.py index cad4802..8169f85 100644 --- a/services/api/app/models/admin.py +++ b/services/api/app/models/admin.py @@ -35,6 +35,9 @@ class ParseSource(Base): docs_added: Mapped[int] = mapped_column(default=0) # Последний (пока идёт — текущий) прогон, см. ParseRun last_run_id: Mapped[int | None] = mapped_column(nullable=True) + # Позиция продолжения для массовых источников: номер статьи в дампе + # Википедии или токен страницы бакета PMC + resume_token: Mapped[str | None] = mapped_column(Text, nullable=True) created_at: Mapped[datetime] = mapped_column(server_default=func.now()) def __repr__(self) -> str: diff --git a/services/api/app/schemas/admin.py b/services/api/app/schemas/admin.py index ff433f0..0ef34ed 100644 --- a/services/api/app/schemas/admin.py +++ b/services/api/app/schemas/admin.py @@ -120,7 +120,9 @@ class ParseRunDetail(ParseRunResponse): # ─── Источники парсинга ─────────────────────────────────────────────────────── class ParseSourceCreate(BaseModel): - source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") + source_type: str = Field( + description="openalex / cyberleninka / arxiv / pmc / wikipedia_ru / pmc_bulk" + ) name: str = Field(min_length=1, max_length=255) query: str | None = None lang: str | None = None @@ -164,7 +166,9 @@ class ParseSourceResponse(BaseModel): class ParseSourceBulkCreate(BaseModel): """Пакетное добавление: один тип/лимит, много запросов (по строке на тему).""" - source_type: str = Field(description="openalex / cyberleninka / arxiv / pmc") + source_type: str = Field( + description="openalex / cyberleninka / arxiv / pmc / wikipedia_ru / pmc_bulk" + ) queries: list[str] = Field(min_length=1, max_length=500) name_prefix: str = "" lang: str | None = None diff --git a/services/frontend/src/pages/admin/Sources.tsx b/services/frontend/src/pages/admin/Sources.tsx index 6f72564..faa818c 100644 --- a/services/frontend/src/pages/admin/Sources.tsx +++ b/services/frontend/src/pages/admin/Sources.tsx @@ -32,7 +32,11 @@ const SOURCE_TYPES = [ { value: 'openalex', label: 'OpenAlex' }, { value: 'cyberleninka', label: 'КиберЛенинка' }, { value: 'arxiv', label: 'arXiv' }, - { value: 'pmc', label: 'PubMed Central' }, + { value: 'pmc', label: 'PubMed Central (по теме)' }, + // Массовые: читают дамп/бакет потоком и продолжают с места остановки, + // поэтому запускаются повторно до исчерпания источника + { value: 'wikipedia_ru', label: 'Википедия ru (дамп, массово)' }, + { value: 'pmc_bulk', label: 'PubMed Central (бакет, массово)' }, ]; const STATUS_COLORS: Record = { diff --git a/services/worker-indexer/app/bulk_writer.py b/services/worker-indexer/app/bulk_writer.py new file mode 100644 index 0000000..8812422 --- /dev/null +++ b/services/worker-indexer/app/bulk_writer.py @@ -0,0 +1,132 @@ +"""Пакетная запись документов корпуса — путь для массовых источников. + +Обычный `index.add_document` пишет по одному документу через ORM: удобно для +парсеров, отдающих сотни статей, но на миллионах это тупик. Замер на проде: +построчная вставка отпечатков — 6.7 тыс. строк/с, COPY — 21 тыс. строк/с, а +миллион статей это ~2 млрд отпечатков. Разница между сутками и неделями. + +Поэтому массовые источники (дампы Википедии, бакет PMC) идут сюда: документы +вставляются одной командой с ON CONFLICT, отпечатки — через COPY. + +Elasticsearch и эмбеддинги здесь намеренно не трогаются: L1 и L2 начинают +работать сразу, а векторы для L3 досчитываются отдельно +(`scripts/ops/reembed_missing.py`) — иначе заливка упирается в скорость +сервера эмбеддингов и тормозит на порядок. +""" + +import contextlib +import io +import logging +import time +from typing import Any + +import psycopg2 + +from app.algorithms.winnowing import sample_evenly, winnow_ordered +from app.config import settings + +logger = logging.getLogger(__name__) + +# Сколько попыток пережить обрыв соединения с базой. Массовая заливка идёт +# часами, и разрыв сети за это время — норма, а не исключение. +DB_RETRIES = 5 + + +def connect(): + """Отдельное подключение psycopg2: COPY недоступен через ORM-сессию.""" + return psycopg2.connect( + host=settings.POSTGRES_HOST, + port=settings.POSTGRES_PORT, + dbname=settings.POSTGRES_DB, + user=settings.POSTGRES_USER, + password=settings.POSTGRES_PASSWORD, + ) + + +def write_batch(conn, docs: list[dict[str, Any]], fp_limit: int) -> tuple[Any, dict[str, int]]: + """Записать пачку документов с отпечатками, пережив обрыв соединения. + + Args: + conn: активное подключение psycopg2 (может быть заменено при обрыве) + docs: документы в унифицированном формате парсеров; нужны ext_id, + source, title и text — остальное необязательно + fp_limit: максимум отпечатков на документ + + Returns: + (соединение, счётчики) — соединение может оказаться новым, если + пришлось переподключаться; счётчики: added / duplicates / fingerprints + """ + if not docs: + return conn, {"added": 0, "duplicates": 0, "fingerprints": 0} + + for attempt in range(1, DB_RETRIES + 1): + try: + return conn, _write_once(conn, docs, fp_limit) + except psycopg2.OperationalError as e: + wait = min(60, 5 * attempt) + logger.warning( + "bulk_writer: база недоступна (%s), повтор через %sс [%s/%s]", + str(e).strip()[:80], wait, attempt, DB_RETRIES, + ) + time.sleep(wait) + with contextlib.suppress(Exception): + conn.close() + try: + conn = connect() + except Exception: + continue + + raise RuntimeError(f"не удалось записать пачку после {DB_RETRIES} попыток") + + +def _write_once(conn, docs: list[dict[str, Any]], fp_limit: int) -> dict[str, int]: + """Одна попытка записи; счётчики возвращаются только при успехе.""" + added = duplicates = fingerprints = 0 + + with conn.cursor() as cur: + cur.executemany( + """INSERT INTO documents (ext_id, source, title, doi, year, lang, journal, + abstract, url, authors, indexed_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,'[]'::json, now()) + ON CONFLICT (ext_id) DO NOTHING""", + [ + ( + d["ext_id"], d["source"], (d.get("title") or "")[:1000], d.get("doi"), + d.get("year"), d.get("lang"), (d.get("journal") or None), + (d.get("text") or "")[:2000], d.get("url"), + ) + 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 = {ext: doc_id for doc_id, ext in cur.fetchall()} + + 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(): + duplicates += 1 + continue + + hashes = sample_evenly(winnow_ordered(d.get("text") or ""), fp_limit) + if not hashes: + continue + for pos, h in enumerate(hashes): + buf.write(f"{doc_id}\t{h}\t{pos}\n") + fingerprints += len(hashes) + added += 1 + + buf.seek(0) + if added: + cur.copy_from(buf, "fingerprints", columns=("doc_id", "hash_value", "position")) + + conn.commit() + return {"added": added, "duplicates": duplicates, "fingerprints": fingerprints} diff --git a/services/worker-indexer/app/config.py b/services/worker-indexer/app/config.py index 5fa2031..adea986 100644 --- a/services/worker-indexer/app/config.py +++ b/services/worker-indexer/app/config.py @@ -59,6 +59,12 @@ class Settings(BaseSettings): FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов + # Сколько документов массового источника пишется одной транзакцией. + # Больше — меньше накладных расходов, но длиннее транзакция: 500 статей по + # 2000 отпечатков это миллион строк за раз, и на таком объёме соединение + # обрывалось. 150 — компромисс, проверенный на заливке Википедии. + BULK_WRITE_BATCH: int = 150 + # Бюджет времени одного прогона index.run_parser. RabbitMQ рвёт канал # consumer'а, не сделавшего ack за consumer_timeout (по умолчанию 1800с): # воркер падает, сообщение передоставляется, таск начинается заново — diff --git a/services/worker-indexer/app/models/__init__.py b/services/worker-indexer/app/models/__init__.py index 1143a1a..6ae9201 100644 --- a/services/worker-indexer/app/models/__init__.py +++ b/services/worker-indexer/app/models/__init__.py @@ -78,6 +78,7 @@ class ParseSource(Base): last_run_at: Mapped[datetime | None] = mapped_column(nullable=True) docs_added: Mapped[int] = mapped_column(default=0) last_run_id: Mapped[int | None] = mapped_column(nullable=True) + resume_token: Mapped[str | None] = mapped_column(Text, nullable=True) created_at: Mapped[datetime] = mapped_column(server_default=func.now()) diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index 3b0f379..95ab943 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -1,5 +1,6 @@ """Celery задачи индексации документов и проверки плагиата (уровни 1-2).""" +import contextlib import io from datetime import UTC, datetime from pathlib import Path @@ -560,12 +561,180 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A "year_from": cfg.get("year_from"), "year_to": cfg.get("year_to"), } + elif source_type == "wikipedia_ru": + # Массовый источник: дамп читается с позиции прошлого прогона + from wikipedia_ru import WikipediaRuParser as P + kwargs = { + "limit": limit, + "dump_path": query or None, # query = путь к дампу, если задан + "skip": int(cfg.get("resume_token") or 0), + } + elif source_type == "pmc_bulk": + from pmc_bulk import PMCBulkParser as P + kwargs = {"limit": limit, "resume_token": cfg.get("resume_token")} else: raise ValueError(f"неизвестный тип источника: {source_type}") return P(), kwargs +def _ingest_one_by_one(parser: Any, raw_docs: Any, prog: Any, run_id: int) -> tuple[bool, bool]: + """Обычный путь: документ за документом через add_document. + + Подходит источникам, отдающим сотни статей: работает ORM-логика с + дедупликацией, индексацией в Elasticsearch и обогащением полным текстом. + + Returns: + (отменён, исчерпан бюджет времени) + """ + cancelled = exhausted = False + + prog.stage = "index" + prog.count_fetched(len(raw_docs)) + prog.log("info", f"получено {len(raw_docs)} документов, индексация") + _write_run(run_id, **prog.snapshot()) + + # Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu + # кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений, + # а не на каждый документ. + embed_batch: list[int] = [] + batch_size = settings.EMBED_BATCH_SIZE + + def _flush_embed() -> None: + if embed_batch: + celery_app.send_task( + "gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu" + ) + embed_batch.clear() + + for raw in raw_docs: + if not cancelled and prog.over_budget: + exhausted = True + prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации") + if cancelled or exhausted: + break + try: + doc = parser.transform(raw) + if not (doc and doc.get("title") and doc.get("ext_id")): + prog.count_result("skipped") + continue + result = add_document(doc, dispatch_embed=False) + prog.count_result(result.get("status", "error")) + if result.get("status") == "indexed": + embed_batch.append(result["doc_id"]) + if len(embed_batch) >= batch_size: + _flush_embed() + except Exception as e: + prog.count_result("error") + prog.log("warning", f"документ пропущен: {e}") + logger.warning(f"run_parser: ошибка документа: {e}") + + if prog.should_flush(): + cancelled = _write_run(run_id, **prog.snapshot()) + prog.mark_flushed() + if cancelled: + prog.log("warning", "отмена по запросу из админки") + + _flush_embed() + return cancelled, exhausted + + +def _ingest_bulk( + parser: Any, raw_docs: Any, prog: Any, run_id: int, source_id: int +) -> tuple[bool, bool]: + """Массовый путь: статьи приходят потоком и пишутся пачками через COPY. + + Отличия от обычного пути и почему они нужны: + - парсер отдаёт генератор, поэтому счётчик «получено» растёт по ходу дела, + а не известен заранее; + - запись идёт пакетно (bulk_writer), иначе миллионы отпечатков занимают + недели вместо суток; + - позиция продолжения сохраняется в источнике, чтобы следующий прогон + начинал с места остановки, а не перечитывал дамп с начала. + + Эмбеддинги здесь не диспатчатся: они считаются заметно медленнее заливки и + делали бы её узким местом. Векторы досчитываются отдельно — + `scripts/ops/reembed_missing.py`. + + Returns: + (отменён, исчерпан бюджет времени) + """ + from app.bulk_writer import connect, write_batch + from app.models import ParseSource + + cancelled = exhausted = False + prog.stage = "index" + prog.log("info", "массовый источник: запись пачками") + _write_run(run_id, **prog.snapshot()) + + conn = connect() + batch: list[dict[str, Any]] = [] + resume_token: str | None = None + batch_size = settings.BULK_WRITE_BATCH + + def save_resume() -> None: + """Запомнить позицию в источнике — с неё продолжит следующий прогон.""" + if resume_token is None: + return + with db_session() as session: + src = session.get(ParseSource, source_id) + if src: + src.resume_token = str(resume_token) + session.commit() + + def flush() -> bool: + """Записать накопленное; True — попросили остановиться.""" + nonlocal conn, batch + if not batch: + return False + conn, counts = write_batch(conn, batch, settings.MAX_FINGERPRINTS_PER_DOC) + for _ in range(counts["added"]): + prog.count_result("indexed") + for _ in range(counts["duplicates"]): + prog.count_result("duplicate") + batch = [] + save_resume() + stop = _write_run(run_id, **prog.snapshot()) + prog.mark_flushed() + return stop + + try: + for raw in raw_docs: + if prog.over_budget: + exhausted = True + prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан") + break + try: + doc = parser.transform(raw) + except Exception as e: + prog.count_result("error") + logger.warning(f"run_parser bulk: ошибка документа: {e}") + continue + + if not (doc and doc.get("ext_id") and doc.get("text")): + prog.count_result("skipped") + continue + + # Позиция продолжения: у Википедии — номер статьи в дампе, + # у PMC — токен страницы бакета + resume_token = doc.get("dump_position") or doc.get("resume_token") or resume_token + batch.append(doc) + prog.count_fetched(prog.fetched + 1) + + if len(batch) >= batch_size and flush(): + cancelled = True + prog.log("warning", "отмена по запросу из админки") + break + + if not cancelled and flush(): + cancelled = True + finally: + with contextlib.suppress(Exception): + conn.close() + + return cancelled, exhausted + + def _write_run(run_id: int, **fields: Any) -> bool: """Записать прогресс прогона и вернуть True, если запрошена отмена. @@ -622,6 +791,8 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any "year_from": src.year_from, "year_to": src.year_to, "limit": src.limit, + # Массовые источники продолжают с сохранённой позиции + "resume_token": src.resume_token, } src.last_status = "running" src.last_error = None @@ -676,53 +847,13 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any # fetch+transform без записи в JSONL — работаем in-memory raw_docs = parser.fetch(**fetch_kwargs, progress_cb=on_fetch_progress) - prog.stage = "index" - prog.count_fetched(len(raw_docs)) - prog.log("info", f"получено {len(raw_docs)} документов, индексация") - _write_run(run_id, **prog.snapshot()) - - # Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu - # кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений, - # а не на каждый документ. - embed_batch: list[int] = [] - batch_size = settings.EMBED_BATCH_SIZE - - def _flush_embed() -> None: - if embed_batch: - celery_app.send_task( - "gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu" - ) - embed_batch.clear() - - for raw in raw_docs: - if not cancelled and prog.over_budget: - exhausted = True - prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации") - if cancelled or exhausted: - break - try: - doc = parser.transform(raw) - if not (doc and doc.get("title") and doc.get("ext_id")): - prog.count_result("skipped") - continue - result = add_document(doc, dispatch_embed=False) - prog.count_result(result.get("status", "error")) - if result.get("status") == "indexed": - embed_batch.append(result["doc_id"]) - if len(embed_batch) >= batch_size: - _flush_embed() - except Exception as e: - prog.count_result("error") - prog.log("warning", f"документ пропущен: {e}") - logger.warning(f"run_parser: ошибка документа: {e}") - - if prog.should_flush(): - cancelled = _write_run(run_id, **prog.snapshot()) - prog.mark_flushed() - if cancelled: - prog.log("warning", "отмена по запросу из админки") - - _flush_embed() + if getattr(parser, "bulk", False): + # Массовые источники (дамп Википедии, бакет PMC) отдают генератор: + # миллионы статей нельзя ни держать в памяти, ни писать по одной + # через ORM — для них отдельный путь с пакетной записью + cancelled, exhausted = _ingest_bulk(parser, raw_docs, prog, run_id, source_id) + else: + cancelled, exhausted = _ingest_one_by_one(parser, raw_docs, prog, run_id) except Exception as e: error_msg = str(e)[:500] prog.log("error", f"прогон упал: {error_msg}")