feat(admin): честная длительность прогона + добор эмбеддингов
Разбор итогов массовой заливки показал в отладке «среднюю длительность прогона» в 2.8 часа там, где парсинг занимал 50 секунд: started_at пишется в момент постановки в очередь, а очередь из 173 источников разбирается часами. Теперь момент реального старта пишется отдельно (run_started_at, миграция 006), а схема отдаёт обе величины — сколько ждал очереди и сколько работал. Плюс scripts/ops/reembed_missing.py: документ попадает в корпус сразу, а вектор для L3 считает отдельная задача; когда worker-gpu или Ollama недоступны, эти задачи теряются и документ остаётся невидимым для семантического поиска. Скрипт находит faiss_id IS NULL и переотправляет пачками (dry-run по умолчанию) — сейчас таких 23 101 из 177 147. Документация: актуальные цифры корпуса, дубли при повторном прогоне, лимит OpenAlex, и главное — гипервизор .254 зафиксирован в DR-HA как самая широкая единая точка отказа (брокер, эмбеддинги, секреты и прокси на одном железе; подтверждено аварией 28.08). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -81,7 +81,10 @@
|
||||
`last_run_id` — ссылка на последний прогон.
|
||||
- **parse_runs** — прогоны заливки: стадия, счётчики (`target/fetched/processed/
|
||||
added/duplicates/skipped/failed`), `cancel_requested`, `heartbeat_at`, журнал
|
||||
событий (JSON). Из них админка рисует шкалу загрузки, см. §12.
|
||||
событий (JSON). Тайминги раздельные: `started_at` — постановка в очередь,
|
||||
`run_started_at` — реальный старт работы воркером (при массовом запуске между
|
||||
ними часы ожидания), `finished_at` — конец. Из них админка рисует шкалу
|
||||
загрузки, см. §12.
|
||||
- **staged_works** — пользовательские загрузки на модерацию перед добавлением в корпус.
|
||||
- **admin_sessions** — одноразовые коды входа в админку.
|
||||
|
||||
@@ -226,6 +229,12 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica
|
||||
`VECTOR_BACKEND`).
|
||||
- **Бюджет времени прогона** (`PARSER_TIME_BUDGET_S`, по умолчанию 1500с) —
|
||||
защита от краш-лупа по `consumer_timeout` RabbitMQ, см. [DR-HA.md](DR-HA.md) §6.
|
||||
- **Добор эмбеддингов** — [`scripts/ops/reembed_missing.py`](../scripts/ops/reembed_missing.py):
|
||||
документ попадает в корпус сразу, а вектор для L3 считает отдельная задача
|
||||
`gpu.embed_documents`; если worker-gpu или Ollama были недоступны, эти задачи
|
||||
теряются и документ остаётся невидимым для семантического поиска. Скрипт
|
||||
находит `faiss_id IS NULL` и переотправляет задачи пачками (dry-run по
|
||||
умолчанию). Покрытие видно в панели отладки.
|
||||
|
||||
## 13. Безопасность
|
||||
|
||||
|
||||
@@ -55,6 +55,17 @@ Redis у нас — кэш/rate-limits/LSH-индекс (префикс `antipla
|
||||
| Векторный индекс | было | `VECTOR_BACKEND=qdrant` снимает (см. README) ✅ |
|
||||
| RabbitMQ .82 | да | мониторинг ловит падение ✅; кластер — по потребности |
|
||||
| app/ES 1.32 | да | воркеры горизонтальны; ES single-node (для BM25 не критично) |
|
||||
| Хост Proxmox .254 | **да, широкий** | митигации нет — см. ниже |
|
||||
|
||||
**Гипервизор — самая широкая единая точка отказа.** На хосте `192.168.20.254`
|
||||
одновременно живут RabbitMQ (.82), embedding-gpu (.109), Infisical (.111) и
|
||||
CT 102 (.253) — фронтенд и реверс-прокси. Его падение снимает сразу: приём и
|
||||
обработку задач (нет брокера), эмбеддинги (нет Ollama), деплой (нет Infisical и
|
||||
Gitea) и весь публичный доступ (нет прокси) — при том что API, PostgreSQL и
|
||||
MinIO продолжают работать. Проверено на практике 2026-08-28: хост перестал
|
||||
отвечать даже на ARP, всё перечисленное отвалилось разом, данные не пострадали.
|
||||
Разнести хотя бы прокси/брокер по разным физическим хостам — самая дешёвая
|
||||
мера; пока её нет, восстановление требует физического доступа к железу.
|
||||
|
||||
## 6. Известный операционный риск — RabbitMQ `consumer_timeout` vs долгие таски
|
||||
|
||||
|
||||
@@ -1,13 +1,18 @@
|
||||
# Наполнение корпуса — runbook
|
||||
|
||||
## Текущее состояние (на 2026-08-27)
|
||||
## Текущее состояние (на 2026-08-28)
|
||||
|
||||
- **165 480 документов**, русский теперь большинство: `ru` 97 507, `en` 67 646,
|
||||
- **177 147 документов**, русский большинство: `ru` 99 884, `en` 76 936,
|
||||
остальные языки — единицы/десятки. Проблема «не с чем сравнивать русские
|
||||
работы» из более ранней версии этого документа закрыта.
|
||||
- По источникам: CyberLeninka 97 506, OpenAlex 49 218, PMC 10 451, arXiv 8 304,
|
||||
- По источникам: CyberLeninka 99 883, OpenAlex 49 276, PMC 14 910, arXiv 13 077,
|
||||
`user_submission` (проверенные пользователями работы, не источник для сравнения
|
||||
сами с собой — см. ARCHITECTURE.md §7) 1.
|
||||
- Повторный прогон по уже залитым источникам даёт почти одни дубли (типично
|
||||
1490 из 1500 на источник): прирост дают только новые публикации. Реальный
|
||||
рост корпуса — поднятый `limit` или новые темы, а не повторный запуск.
|
||||
- OpenAlex при массовом запуске упирается в лимит вежливого пула (10 req/s
|
||||
на mailto, общий для всех воркеров) — пауза между страницами поднята до 1с.
|
||||
- Добавлен 4-й парсер — **PMC** (PubMed Central, `scripts/parsers/pmc.py`),
|
||||
англоязычные научные статьи открытого доступа.
|
||||
- Массовое расширение по дисциплинам теперь двумя сидерами: `seed_ru_sources.py`
|
||||
@@ -61,6 +66,11 @@ python scripts/seed_ru_sources.py --apply --limit 500
|
||||
- **Страница «Отладка»** — очереди, воркеры, покрытие эмбеддингами, зависшие и
|
||||
упавшие прогоны (см. ARCHITECTURE.md §12).
|
||||
|
||||
После большой заливки стоит свериться с покрытием L3 (панель отладки, строка
|
||||
«без вектора»): эмбеддинги считаются отдельной задачей на worker-gpu, и если он
|
||||
или Ollama были недоступны, документы останутся без вектора. Догнать —
|
||||
`scripts/ops/reembed_missing.py` (dry-run по умолчанию, `--apply` отправляет).
|
||||
|
||||
Прогон со статусом `partial` — это не ошибка: сработал бюджет времени
|
||||
(`PARSER_TIME_BUDGET_S`, 1500с), заливка остановилась раньше `consumer_timeout`
|
||||
RabbitMQ. Остаток добирается повторным запуском источника.
|
||||
@@ -82,6 +92,6 @@ SELECT source, count(*) FROM documents WHERE source='cyberleninka'; -- > 0
|
||||
заголовки с меткой ru).
|
||||
- Для миллионов — **bulk** (снапшот OpenAlex на S3), а не постраничный API.
|
||||
- На масштабе обязателен `VECTOR_BACKEND=qdrant` (FAISS flat не тянет), а таблица
|
||||
`fingerprints` (уже ~88M строк на 165K доков — партиционирование стоит планировать
|
||||
`fingerprints` (уже ~113M строк на 177K доков — партиционирование стоит планировать
|
||||
заранее, не постфактум) потребует партиционирования. См.
|
||||
[ARCHITECTURE.md](ARCHITECTURE.md) и [DR-HA.md](DR-HA.md).
|
||||
|
||||
140
scripts/ops/reembed_missing.py
Executable file
140
scripts/ops/reembed_missing.py
Executable file
@@ -0,0 +1,140 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Догнать эмбеддинги для документов, у которых их нет (faiss_id IS NULL).
|
||||
|
||||
Зачем: документ попадает в корпус сразу (метаданные, fingerprints L1, MinHash L2),
|
||||
а вектор для L3 считает отдельная задача `gpu.embed_documents`. Если worker-gpu
|
||||
или Ollama были недоступны в момент заливки, эти задачи теряются — документ
|
||||
остаётся в базе, но семантический поиск его не видит. Скрипт находит такие
|
||||
документы и переотправляет задачи пачками.
|
||||
|
||||
Идемпотентно: повторный запуск возьмёт только оставшиеся без faiss_id.
|
||||
|
||||
ВАЖНО, чего скрипт НЕ делает: документы с непустым `faiss_id`, которого нет в
|
||||
самом FAISS-индексе (наследие смены модели эмбеддингов — размерность сменилась,
|
||||
индекс пересобран с нуля, а ссылки в БД остались), сюда не попадут. Их сначала
|
||||
надо выявить сверкой с индексом и обнулить faiss_id — это отдельная операция.
|
||||
|
||||
Запуск (нужны psycopg2 + celery — они есть в образе воркера; сам репозиторий
|
||||
внутрь контейнера не смонтирован, поэтому передаём скрипт через stdin):
|
||||
cd /home/user/anti-plagiarism
|
||||
C="docker compose -f docker-compose.prod.yml exec -T worker-indexer python -"
|
||||
$C < scripts/ops/reembed_missing.py # dry-run, только считает
|
||||
$C --apply < scripts/ops/reembed_missing.py # отправить задачи
|
||||
$C --apply --limit 1000 < scripts/ops/reembed_missing.py # пробная порция
|
||||
|
||||
Креды берутся из окружения (в контейнере они уже есть из env_file), а при
|
||||
запуске файлом — из .env репозитория: POSTGRES_*, RABBITMQ_URL.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import os
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
DEFAULT_BATCH = 64 # столько же, сколько EMBED_BATCH_SIZE у индексера
|
||||
|
||||
|
||||
def _load_env() -> None:
|
||||
"""Подтянуть переменные из .env репозитория, если файл есть.
|
||||
|
||||
При запуске через stdin (`docker exec ... python -`) __file__ не определён —
|
||||
это штатный способ запуска здесь, и переменные в контейнере уже есть.
|
||||
"""
|
||||
try:
|
||||
env = Path(__file__).resolve().parent.parent.parent / ".env"
|
||||
except NameError:
|
||||
return
|
||||
if not env.exists():
|
||||
return
|
||||
for line in env.read_text(encoding="utf-8").splitlines():
|
||||
line = line.strip()
|
||||
if not line or line.startswith("#") or "=" not in line:
|
||||
continue
|
||||
key, _, val = line.partition("=")
|
||||
os.environ.setdefault(key.strip(), val.strip())
|
||||
|
||||
|
||||
def _pg_dsn() -> dict:
|
||||
return {
|
||||
"host": os.environ.get("POSTGRES_HOST", "localhost"),
|
||||
"port": int(os.environ.get("POSTGRES_PORT", "5432")),
|
||||
"dbname": os.environ.get("POSTGRES_DB", "antiplagiator"),
|
||||
"user": os.environ.get("POSTGRES_USER", "antiplagiator"),
|
||||
"password": os.environ.get("POSTGRES_PASSWORD", ""),
|
||||
}
|
||||
|
||||
|
||||
def chunked(items: list[int], size: int) -> list[list[int]]:
|
||||
"""Разбить список id на пачки по size элементов."""
|
||||
return [items[i : i + size] for i in range(0, len(items), size)]
|
||||
|
||||
|
||||
def main() -> None:
|
||||
ap = argparse.ArgumentParser(
|
||||
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
|
||||
)
|
||||
ap.add_argument("--apply", action="store_true", help="реально отправить задачи (иначе dry-run)")
|
||||
ap.add_argument("--batch", type=int, default=DEFAULT_BATCH, help=f"документов в пачке (по умолчанию {DEFAULT_BATCH})")
|
||||
ap.add_argument("--limit", type=int, default=0, help="взять не больше N документов (0 = все)")
|
||||
args = ap.parse_args()
|
||||
|
||||
_load_env()
|
||||
|
||||
try:
|
||||
import psycopg2
|
||||
except ImportError:
|
||||
sys.exit("Нужен psycopg2 (есть в образах воркеров). Запускайте в контейнере.")
|
||||
|
||||
dsn = _pg_dsn()
|
||||
print(f"PostgreSQL {dsn['host']}:{dsn['port']}/{dsn['dbname']}")
|
||||
conn = psycopg2.connect(**dsn)
|
||||
with conn.cursor() as cur:
|
||||
sql = "SELECT id FROM documents WHERE faiss_id IS NULL ORDER BY id"
|
||||
if args.limit:
|
||||
sql += f" LIMIT {int(args.limit)}"
|
||||
cur.execute(sql)
|
||||
doc_ids = [r[0] for r in cur.fetchall()]
|
||||
|
||||
cur.execute("SELECT count(*) FROM documents")
|
||||
total = cur.fetchone()[0]
|
||||
conn.close()
|
||||
|
||||
batches = chunked(doc_ids, args.batch)
|
||||
print(f"Документов всего: {total}")
|
||||
print(f"Без эмбеддинга: {len(doc_ids)} → пачек по {args.batch}: {len(batches)}")
|
||||
if not doc_ids:
|
||||
print("Нечего досчитывать.")
|
||||
return
|
||||
|
||||
if not args.apply:
|
||||
head = ", ".join(str(i) for i in doc_ids[:10])
|
||||
print(f"\n[dry-run] первые id: {head}{' …' if len(doc_ids) > 10 else ''}")
|
||||
print("Запустите с --apply, чтобы отправить gpu.embed_documents в queue.gpu.")
|
||||
return
|
||||
|
||||
try:
|
||||
from celery import Celery
|
||||
except ImportError:
|
||||
sys.exit("Нужен celery (есть в образах воркеров). Запускайте в контейнере.")
|
||||
|
||||
broker = os.environ.get("RABBITMQ_URL")
|
||||
if not broker:
|
||||
sys.exit("Не задан RABBITMQ_URL — нечем отправлять задачи.")
|
||||
|
||||
app = Celery("reembed", broker=broker)
|
||||
sent = 0
|
||||
try:
|
||||
for batch in batches:
|
||||
app.send_task("gpu.embed_documents", args=[batch], queue="queue.gpu")
|
||||
sent += 1
|
||||
if sent % 50 == 0:
|
||||
print(f" отправлено пачек: {sent}/{len(batches)}")
|
||||
except Exception as e:
|
||||
sys.exit(f"\nБрокер недоступен после {sent} пачек: {e}")
|
||||
|
||||
print(f"\nОтправлено пачек: {sent} ({len(doc_ids)} документов).")
|
||||
print("Ход выполнения — админка → «Отладка» (очередь queue.gpu и покрытие эмбеддингами).")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
31
services/api/alembic/versions/006_parse_run_started_at.py
Normal file
31
services/api/alembic/versions/006_parse_run_started_at.py
Normal file
@@ -0,0 +1,31 @@
|
||||
"""Отдельная отметка старта выполнения прогона парсинга.
|
||||
|
||||
Revision ID: 006
|
||||
Revises: 005
|
||||
Create Date: 2026-08-28
|
||||
|
||||
`started_at` пишется в момент СОЗДАНИЯ строки, то есть постановки в очередь.
|
||||
При массовом запуске очередь разбирается часами, и «длительность» прогона по
|
||||
двум таймстампам показывала 2.8 часа там, где сам парсинг занял 50 секунд —
|
||||
для отладки это дезинформация. Момент, когда воркер реально взял задачу,
|
||||
пишем отдельно; разница со `started_at` — это ожидание в очереди.
|
||||
"""
|
||||
|
||||
import sqlalchemy as sa
|
||||
|
||||
from alembic import op
|
||||
|
||||
revision = "006"
|
||||
down_revision = "005"
|
||||
branch_labels = None
|
||||
depends_on = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
# Для уже прошедших прогонов остаётся NULL: подставлять им started_at
|
||||
# значило бы выдать время ожидания за время работы.
|
||||
op.add_column("parse_runs", sa.Column("run_started_at", sa.DateTime(), nullable=True))
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_column("parse_runs", "run_started_at")
|
||||
@@ -75,7 +75,10 @@ class ParseRun(Base):
|
||||
error: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
# [{"ts": ISO8601, "level": "info|warning|error", "msg": str}]
|
||||
log: Mapped[list | None] = mapped_column(JSON, default=list)
|
||||
# Постановка в очередь (строку создаёт API) и реальный старт выполнения:
|
||||
# между ними при массовом запуске проходят часы, мешать их нельзя
|
||||
started_at: Mapped[datetime] = mapped_column(server_default=func.now(), index=True)
|
||||
run_started_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
# Последний признак жизни: по нему видно зависший прогон (running, но тишина)
|
||||
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
|
||||
@@ -80,7 +80,9 @@ class ParseRunResponse(BaseModel):
|
||||
failed: int
|
||||
cancel_requested: bool = False
|
||||
error: str | None = None
|
||||
# started_at — постановка в очередь, run_started_at — реальный старт работы
|
||||
started_at: datetime
|
||||
run_started_at: datetime | None = None
|
||||
heartbeat_at: datetime | None = None
|
||||
finished_at: datetime | None = None
|
||||
|
||||
@@ -91,6 +93,26 @@ class ParseRunResponse(BaseModel):
|
||||
def percent(self) -> float:
|
||||
return run_percent(self.status, self.stage, self.target, self.fetched, self.processed)
|
||||
|
||||
@computed_field # type: ignore[prop-decorator]
|
||||
@property
|
||||
def queued_s(self) -> float | None:
|
||||
"""Сколько прогон ждал своей очереди, сек."""
|
||||
if self.run_started_at is None:
|
||||
return None
|
||||
return round((self.run_started_at - self.started_at).total_seconds(), 1)
|
||||
|
||||
@computed_field # type: ignore[prop-decorator]
|
||||
@property
|
||||
def duration_s(self) -> float | None:
|
||||
"""Сколько прогон реально работал, сек (None — ещё идёт или не начинался).
|
||||
|
||||
Прогоны до появления run_started_at (миграция 006) остаются без
|
||||
длительности: у них известен только момент постановки в очередь.
|
||||
"""
|
||||
if self.run_started_at is None or self.finished_at is None:
|
||||
return None
|
||||
return round((self.finished_at - self.run_started_at).total_seconds(), 1)
|
||||
|
||||
|
||||
class ParseRunDetail(ParseRunResponse):
|
||||
log: list[dict[str, Any]] = Field(default_factory=list)
|
||||
|
||||
@@ -13,8 +13,9 @@ interface Run {
|
||||
target: number; fetched: number; processed: number;
|
||||
added: number; duplicates: number; skipped: number; failed: number;
|
||||
cancel_requested: boolean; error: string | null;
|
||||
started_at: string; heartbeat_at: string | null; finished_at: string | null;
|
||||
percent: number;
|
||||
started_at: string; run_started_at: string | null;
|
||||
heartbeat_at: string | null; finished_at: string | null;
|
||||
percent: number; queued_s: number | null; duration_s: number | null;
|
||||
}
|
||||
|
||||
interface LogEntry { ts: string; elapsed: number; level: string; msg: string }
|
||||
@@ -360,6 +361,14 @@ export function Sources() {
|
||||
);
|
||||
}
|
||||
|
||||
/** Секунды → «45с» / «12 мин» / «2 ч 5 мин»: в отладке важен порядок, не точность. */
|
||||
function fmtDuration(seconds: number): string {
|
||||
if (seconds < 90) return `${Math.round(seconds)}с`;
|
||||
const min = Math.round(seconds / 60);
|
||||
if (min < 90) return `${min} мин`;
|
||||
return `${Math.floor(min / 60)} ч ${min % 60} мин`;
|
||||
}
|
||||
|
||||
const LOG_COLORS: Record<string, string> = {
|
||||
error: 'text-red-600',
|
||||
warning: 'text-amber-600',
|
||||
@@ -389,7 +398,9 @@ function RunLog({ runId, live }: { runId?: number; live: boolean }) {
|
||||
<Chip label="дублей" value={String(data.duplicates)} />
|
||||
{data.skipped > 0 && <Chip label="без метаданных" value={String(data.skipped)} />}
|
||||
{data.failed > 0 && <Chip label="ошибок" value={String(data.failed)} />}
|
||||
<Chip label="начат" value={new Date(data.started_at).toLocaleString('ru-RU')} />
|
||||
<Chip label="поставлен в очередь" value={new Date(data.started_at).toLocaleString('ru-RU')} />
|
||||
{data.queued_s != null && <Chip label="ждал очереди" value={fmtDuration(data.queued_s)} />}
|
||||
{data.duration_s != null && <Chip label="работал" value={fmtDuration(data.duration_s)} />}
|
||||
{data.finished_at && <Chip label="завершён" value={new Date(data.finished_at).toLocaleString('ru-RU')} />}
|
||||
</div>
|
||||
|
||||
|
||||
@@ -101,6 +101,7 @@ class ParseRun(Base):
|
||||
error: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
log: Mapped[list | None] = mapped_column(JSON, default=list)
|
||||
started_at: Mapped[datetime] = mapped_column(server_default=func.now())
|
||||
run_started_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(nullable=True)
|
||||
|
||||
|
||||
@@ -625,7 +625,15 @@ def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any
|
||||
prog.stage = "fetch"
|
||||
prog.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}")
|
||||
_write_run(
|
||||
run_id, status="running", celery_task_id=self.request.id, error=None, **prog.snapshot()
|
||||
run_id,
|
||||
status="running",
|
||||
celery_task_id=self.request.id,
|
||||
error=None,
|
||||
# Отдельно от started_at (постановка в очередь): при массовом запуске
|
||||
# между ними часы ожидания, и без этой отметки «длительность прогона»
|
||||
# в отладке показывала очередь, а не работу
|
||||
run_started_at=datetime.now(UTC),
|
||||
**prog.snapshot(),
|
||||
)
|
||||
|
||||
cancelled = False
|
||||
|
||||
Reference in New Issue
Block a user