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>
This commit is contained in:
132
services/worker-indexer/app/bulk_writer.py
Normal file
132
services/worker-indexer/app/bulk_writer.py
Normal file
@@ -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}
|
||||
Reference in New Issue
Block a user