Заливку Википедии и 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>
133 lines
5.9 KiB
Python
133 lines
5.9 KiB
Python
"""Пакетная запись документов корпуса — путь для массовых источников.
|
||
|
||
Обычный `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}
|