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}")