Files
jze9 26de1b9de1
All checks were successful
Deploy / deploy (push) Successful in 2m28s
Deploy / test (push) Successful in 2m49s
refactor(ops): массовая заливка — в общий конвейер вместо отдельных скриптов
Заливку Википедии и 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 <noreply@anthropic.com>
2026-09-05 20:32:53 +05:00

133 lines
5.9 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Пакетная запись документов корпуса — путь для массовых источников.
Обычный `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}