Compare commits

...

3 Commits

Author SHA1 Message Date
jze9
39951d6a0f feat(admin): показывать зависшие проверки в панели отладки
All checks were successful
Deploy / test (push) Successful in 3m44s
Deploy / deploy (push) Successful in 8m39s
Нашлись 4 задачи в статусе processing, висящие с 30 мая: пользователь видит
вечное «обрабатывается», а в системе никаких следов. Теперь задачи без
движения дольше 2 часов попадают в /admin/debug рядом с зависшими прогонами
заливки, с указанием, сколько часов они стоят.

tasks.created_at/updated_at — timestamptz, поэтому сравнение идёт с
осведомлённым о зоне временем: naive utcnow() дал бы TypeError в рантайме.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
jze9
f4092e5331 fix(ops): заливки переживают обрывы базы и сети
За неделю обе массовые заливки умерли молча и не возобновились: Википедия на
20 008 статьях из ~2 млн («server closed the connection»), PMC примерно на
85 тыс. из 100 тыс. (таймаут SSL-рукопожатия). Ни ретраев, ни возобновления.

- запись пачки повторяется с переподключением к базе (5 попыток с паузой);
- обрыв листинга бакета стоит паузы, а не всей заливки;
- у Википедии сохраняется позиция в дампе — после сбоя продолжаем с неё,
  а не проматываем с нуля, полагаясь на дедуп по ext_id;
- батч уменьшен (500 статей × 2000 отпечатков = миллион строк в одной
  транзакции — вероятная причина обрыва);
- обрезка отпечатков переведена на равномерную выборку.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
jze9
476682741e fix(L1): обрезать отпечатки равномерно по тексту, а не произвольно
winnow() возвращает set, поэтому list(fp)[:LIMIT] брал случайное подмножество:
у длинного документа целые куски оставались без отпечатков, и списывание
именно из них не находилось. Обнаружено при разборе того, почему фрагмент
статьи PMC не искался.

Добавлены winnow_ordered() — отпечатки в порядке появления в тексте, и
sample_evenly() — выборка каждого n-го элемента вместо первых N. Применено в
add_document и store_full_text.

7 тестов на главное свойство: выборка растянута по всей длине документа, шаг
ровный, порядок сохранён.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-05 19:13:01 +05:00
6 changed files with 317 additions and 70 deletions

View File

@@ -25,15 +25,15 @@
залитые пропускаются. Прогресс листинга сохраняется в --state-file, поэтому
прерванная заливка продолжается с той же страницы бакета.
Проверено на проде 2026-08-31: 100 статей залито, 183 090 отпечатков (1831 на
статью — реальная глубина, а не аннотация). Темп — 0.3 статьи/с, то есть
миллион в один поток занял бы около 40 суток: скачивание тут не узкое место
(12 статей/с в 12 потоков), упирается в последовательную обработку одного
документа. Для миллионов обработку нужно распараллелить по ядрам воркера —
это следующий шаг, до него скрипт годится для порций в десятки тысяч.
Проверено на проде: 1831 отпечаток на статью — реальная глубина, а не аннотация.
Темп ~2 статьи/с в один поток (0.3 ст/с из первого замера были получены, когда
база одновременно считала эмбеддинги). Скачивание не узкое место (12 статей/с в
12 потоков), упирается в последовательную обработку документа — для миллионов
её нужно распараллелить по ядрам воркера.
"""
import argparse
import contextlib
import io
import json
import os
@@ -93,6 +93,76 @@ def fetch_one(client, art_id: str) -> dict | None:
}
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)
@@ -105,7 +175,6 @@ def main() -> None:
import httpx
import psycopg2
from app.algorithms.winnowing import winnow
from app.config import settings
token = None
@@ -114,19 +183,31 @@ def main() -> None:
token = json.load(fh).get("token")
print("продолжаю листинг с сохранённой позиции")
client = httpx.Client(timeout=60, headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"})
conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
password=settings.POSTGRES_PASSWORD)
client = httpx.Client(timeout=60,
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"})
added = skipped = failed = 0
fp_total = 0
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()
done = 0
while done < args.limit:
want = min(args.batch, args.limit - done)
ids, token = list_articles(client, token, want)
# Сеть у бакета иногда отваливается (прошлый прогон умер на таймауте
# 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
@@ -138,46 +219,10 @@ def main() -> None:
if not docs:
continue
with conn.cursor() as cur:
# Документы: ON CONFLICT — повторный прогон не плодит дубли
rows = [
(d["ext_id"], "pmc", d["title"], d["doi"], d["year"], "en",
d["journal"], d["text"][:2000])
for d in docs
]
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""",
rows,
)
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()
fresh = 0
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
hashes = list(winnow(d["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()
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)
@@ -194,7 +239,8 @@ def main() -> None:
print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, "
f"недоступно {failed}, отпечатков {fp_total:,}")
if added:
print(f"темп: {added / el:.1f} статей/с → миллион за ~{1e6 / max(added / el, 0.01) / 3600:.0f} ч")
print(f"темп: {added / el:.1f} статей/с → "
f"миллион за ~{1e6 / max(added / el, 0.01) / 3600:.0f} ч")
if __name__ == "__main__":

View File

@@ -25,7 +25,10 @@
import argparse
import bz2
import contextlib
import io
import json
import os
import re
import time
@@ -138,29 +141,45 @@ 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=500, 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 winnow
from app.algorithms.winnowing import sample_evenly, winnow_ordered
from app.config import settings
conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
password=settings.POSTGRES_PASSWORD)
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]] = []
def flush(batch):
# Позиция в дампе: прошлый прогон обрывался на 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
if not batch:
return
with conn.cursor() as cur:
cur.executemany(
"""INSERT INTO documents (ext_id, source, title, lang, abstract, url,
@@ -184,7 +203,9 @@ def main() -> None:
if cur.fetchone():
skipped += 1
continue
hashes = list(winnow(text))[: args.fp_per_doc]
# Равномерно по тексту, а не первые 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)
@@ -195,11 +216,42 @@ def main() -> None:
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)

View File

@@ -9,7 +9,7 @@ import io
import logging
import os
import uuid
from datetime import datetime, timedelta
from datetime import UTC, datetime, timedelta
from pathlib import Path
import httpx
@@ -925,6 +925,25 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
)
).scalars().all()
# Зависшая проверка: числится в работе, но воркер давно её не трогал.
# Пользователь видит вечное «обрабатывается» — молча и без следов в логах,
# поэтому такие задачи должны быть видны в отладке (нашлись экземпляры,
# висевшие с мая).
# tasks.created_at/updated_at — timestamptz, поэтому сравниваем с
# осведомлённым о зоне временем (naive utcnow() дал бы TypeError)
task_stale_cutoff = datetime.now(UTC) - timedelta(hours=2)
stale_tasks = (
await db.execute(
select(Task)
.where(
Task.status == "processing",
func.coalesce(Task.updated_at, Task.created_at) < task_stale_cutoff,
)
.order_by(Task.created_at.desc())
.limit(20)
)
).scalars().all()
return {
"generated_at": datetime.utcnow().isoformat(timespec="seconds"),
"celery": celery_state,
@@ -956,6 +975,17 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
}
for t in failed_tasks
],
"stale_tasks": [
{
"public_id": t.public_id,
"type": t.type,
"created_at": t.created_at,
"stuck_hours": round(
(datetime.now(UTC) - (t.updated_at or t.created_at)).total_seconds() / 3600, 1
),
}
for t in stale_tasks
],
# Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на
# поведение воркеров — при разборе «почему не так считает» нужны первыми
"config": {

View File

@@ -85,6 +85,68 @@ def winnow(text: str, k: int = 5, window: int = 4) -> set[int]:
return fingerprint
def winnow_ordered(text: str, k: int = 5, window: int = 4) -> list[int]:
"""То же, что winnow, но с сохранением порядка появления отпечатков в тексте.
Нужно там, где отпечатки приходится обрезать по лимиту: множество не хранит
порядка, и `list(winnow(text))[:limit]` берёт произвольное подмножество —
у длинного документа целые куски остаются без покрытия, и списывание из них
не находится. Со списком в порядке текста обрезку можно делать равномерной
(см. sample_evenly).
Args:
text: Исходный текст
k: Размер k-граммы
window: Размер скользящего окна
Returns:
Отпечатки в порядке появления в тексте, без повторов
"""
tokens = text.lower().split()
if len(tokens) < k:
return []
hashes = [hash_ngram(ng) for ng in get_ngrams(tokens, k)]
if not hashes:
return []
ordered: list[int] = []
seen: set[int] = set()
prev_min_idx = -1
for i in range(len(hashes) - window + 1):
window_hashes = hashes[i : i + window]
min_val = min(window_hashes)
min_idx = i + window_hashes.index(min_val)
if min_idx != prev_min_idx:
if min_val not in seen:
seen.add(min_val)
ordered.append(min_val)
prev_min_idx = min_idx
return ordered
def sample_evenly(items: list[int], limit: int) -> list[int]:
"""Оставить не больше limit элементов, равномерно по всей длине списка.
Берём каждый n-й элемент, а не первые limit штук: обрезка «с начала»
оставила бы без отпечатков весь конец документа.
Args:
items: Отпечатки в порядке текста (winnow_ordered)
limit: Максимум отпечатков; 0 или меньше — не ограничивать
Returns:
Подсписок длиной не больше limit, сохраняющий порядок
"""
if limit <= 0 or len(items) <= limit:
return items
step = len(items) / limit
return [items[int(i * step)] for i in range(limit)]
def jaccard_similarity(fp_a: set[int], fp_b: set[int]) -> float:
"""
Коэффициент Жаккара для двух fingerprint'ов.

View File

@@ -10,7 +10,7 @@ from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError
from app.algorithms.minhash import add_to_lsh, find_similar
from app.algorithms.winnowing import winnow
from app.algorithms.winnowing import sample_evenly, winnow, winnow_ordered
from app.celery_app import celery_app
from app.config import settings
from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status
@@ -284,8 +284,9 @@ def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[
# Вычислить fingerprints
text = doc_data.get("full_text") or doc_data.get("abstract", "") or ""
if text:
fp = winnow(text)
fingerprints_to_add = list(fp)[: settings.MAX_FINGERPRINTS_PER_DOC]
fingerprints_to_add = sample_evenly(
winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC
)
for i, hash_val in enumerate(fingerprints_to_add):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
@@ -375,7 +376,7 @@ def store_full_text(doc_id: int, text: str) -> dict[str, Any]:
content_type="text/plain; charset=utf-8",
)
hashes = list(winnow(text))[: settings.MAX_FINGERPRINTS_PER_DOC]
hashes = sample_evenly(winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC)
with db_session() as session:
doc = session.get(Document, doc_id)

View File

@@ -5,7 +5,9 @@ from app.algorithms.winnowing import (
get_ngrams,
hash_ngram,
jaccard_similarity,
sample_evenly,
winnow,
winnow_ordered,
)
# Достаточно длинный текст, чтобы окно Winnowing реально отработало
@@ -77,3 +79,57 @@ def test_compute_similarity_partial_overlap_is_between():
modified = LONG + " добавлен ещё один совершенно новый хвост предложения здесь"
sim = compute_similarity(LONG, modified)
assert 0.0 < sim < 1.0
# ─── Обрезка отпечатков по лимиту ────────────────────────────────────────────
# Регрессия: winnow() возвращает set, и list(fp)[:limit] брал произвольное
# подмножество — у длинного документа целые куски оставались без покрытия,
# и списывание из них не находилось.
def test_winnow_ordered_matches_winnow_by_content():
"""Тот же набор отпечатков, что и у winnow, только с порядком."""
assert set(winnow_ordered(LONG)) == winnow(LONG)
def test_winnow_ordered_has_no_duplicates():
ordered = winnow_ordered(LONG)
assert len(ordered) == len(set(ordered))
def test_winnow_ordered_follows_text_order():
"""Отпечатки начала текста идут раньше отпечатков продолжения."""
tail = " совершенно другой хвост про выпечку хлеба и закваску в тёплой печи"
ordered = winnow_ordered(LONG + tail)
head_prints = set(winnow_ordered(LONG))
positions = [i for i, h in enumerate(ordered) if h in head_prints]
# Отпечатки первой половины сосредоточены в начале списка, а не разбросаны
assert max(positions) < len(ordered)
assert positions[0] == 0
def test_sample_evenly_keeps_everything_under_limit():
items = [1, 2, 3]
assert sample_evenly(items, 10) == items
assert sample_evenly(items, 0) == items # 0 = без ограничения
def test_sample_evenly_respects_limit():
items = list(range(1000))
assert len(sample_evenly(items, 100)) == 100
def test_sample_evenly_covers_whole_document():
"""Главное свойство: выборка растянута по всей длине, а не обрезана с начала."""
items = list(range(1000))
sampled = sample_evenly(items, 10)
assert sampled[0] == 0
assert sampled[-1] >= 900 # хвост документа тоже покрыт
assert sampled == sorted(sampled) # порядок сохранён
def test_sample_evenly_spreads_uniformly():
items = list(range(100))
sampled = sample_evenly(items, 10)
gaps = [b - a for a, b in zip(sampled, sampled[1:], strict=False)]
assert max(gaps) - min(gaps) <= 1 # шаг ровный