Compare commits
3 Commits
9145de8e54
...
39951d6a0f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
39951d6a0f | ||
|
|
f4092e5331 | ||
|
|
476682741e |
@@ -25,15 +25,15 @@
|
|||||||
залитые пропускаются. Прогресс листинга сохраняется в --state-file, поэтому
|
залитые пропускаются. Прогресс листинга сохраняется в --state-file, поэтому
|
||||||
прерванная заливка продолжается с той же страницы бакета.
|
прерванная заливка продолжается с той же страницы бакета.
|
||||||
|
|
||||||
Проверено на проде 2026-08-31: 100 статей залито, 183 090 отпечатков (1831 на
|
Проверено на проде: 1831 отпечаток на статью — реальная глубина, а не аннотация.
|
||||||
статью — реальная глубина, а не аннотация). Темп — 0.3 статьи/с, то есть
|
Темп ~2 статьи/с в один поток (0.3 ст/с из первого замера были получены, когда
|
||||||
миллион в один поток занял бы около 40 суток: скачивание тут не узкое место
|
база одновременно считала эмбеддинги). Скачивание не узкое место (12 статей/с в
|
||||||
(12 статей/с в 12 потоков), упирается в последовательную обработку одного
|
12 потоков), упирается в последовательную обработку документа — для миллионов
|
||||||
документа. Для миллионов обработку нужно распараллелить по ядрам воркера —
|
её нужно распараллелить по ядрам воркера.
|
||||||
это следующий шаг, до него скрипт годится для порций в десятки тысяч.
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
|
import contextlib
|
||||||
import io
|
import io
|
||||||
import json
|
import json
|
||||||
import os
|
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:
|
def main() -> None:
|
||||||
ap = argparse.ArgumentParser(description=__doc__,
|
ap = argparse.ArgumentParser(description=__doc__,
|
||||||
formatter_class=argparse.RawDescriptionHelpFormatter)
|
formatter_class=argparse.RawDescriptionHelpFormatter)
|
||||||
@@ -105,7 +175,6 @@ def main() -> None:
|
|||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
import psycopg2
|
import psycopg2
|
||||||
from app.algorithms.winnowing import winnow
|
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
|
||||||
token = None
|
token = None
|
||||||
@@ -114,19 +183,31 @@ def main() -> None:
|
|||||||
token = json.load(fh).get("token")
|
token = json.load(fh).get("token")
|
||||||
print("продолжаю листинг с сохранённой позиции")
|
print("продолжаю листинг с сохранённой позиции")
|
||||||
|
|
||||||
client = httpx.Client(timeout=60, headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"})
|
client = httpx.Client(timeout=60,
|
||||||
conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
|
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"})
|
||||||
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
|
|
||||||
password=settings.POSTGRES_PASSWORD)
|
|
||||||
|
|
||||||
added = skipped = failed = 0
|
def connect():
|
||||||
fp_total = 0
|
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()
|
t0 = time.time()
|
||||||
done = 0
|
|
||||||
|
|
||||||
while done < args.limit:
|
while done < args.limit:
|
||||||
want = min(args.batch, args.limit - done)
|
want = min(args.batch, args.limit - done)
|
||||||
|
# Сеть у бакета иногда отваливается (прошлый прогон умер на таймауте
|
||||||
|
# SSL-рукопожатия) — это стоит паузы, а не всей заливки
|
||||||
|
try:
|
||||||
ids, token = list_articles(client, token, want)
|
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:
|
if not ids:
|
||||||
print("бакет закончился")
|
print("бакет закончился")
|
||||||
break
|
break
|
||||||
@@ -138,46 +219,10 @@ def main() -> None:
|
|||||||
if not docs:
|
if not docs:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
with conn.cursor() as cur:
|
conn, n_added, n_skipped, n_fp = write_batch(conn, connect, docs, args.fp_per_doc)
|
||||||
# Документы: ON CONFLICT — повторный прогон не плодит дубли
|
added += n_added
|
||||||
rows = [
|
skipped += n_skipped
|
||||||
(d["ext_id"], "pmc", d["title"], d["doi"], d["year"], "en",
|
fp_total += n_fp
|
||||||
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()
|
|
||||||
|
|
||||||
with open(args.state_file, "w") as fh:
|
with open(args.state_file, "w") as fh:
|
||||||
json.dump({"token": token}, fh)
|
json.dump({"token": token}, fh)
|
||||||
@@ -194,7 +239,8 @@ def main() -> None:
|
|||||||
print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, "
|
print(f"\nГотово за {el / 60:.1f} мин: залито {added}, пропущено {skipped}, "
|
||||||
f"недоступно {failed}, отпечатков {fp_total:,}")
|
f"недоступно {failed}, отпечатков {fp_total:,}")
|
||||||
if added:
|
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__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
@@ -25,7 +25,10 @@
|
|||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
import bz2
|
import bz2
|
||||||
|
import contextlib
|
||||||
import io
|
import io
|
||||||
|
import json
|
||||||
|
import os
|
||||||
import re
|
import re
|
||||||
import time
|
import time
|
||||||
|
|
||||||
@@ -138,29 +141,45 @@ def main() -> None:
|
|||||||
ap = argparse.ArgumentParser(description=__doc__,
|
ap = argparse.ArgumentParser(description=__doc__,
|
||||||
formatter_class=argparse.RawDescriptionHelpFormatter)
|
formatter_class=argparse.RawDescriptionHelpFormatter)
|
||||||
ap.add_argument("--limit", type=int, default=100, help="сколько статей залить")
|
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("--min-chars", type=int, default=2000, help="минимальная длина текста")
|
||||||
ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на статью")
|
ap.add_argument("--fp-per-doc", type=int, default=2000, help="максимум отпечатков на статью")
|
||||||
ap.add_argument("--dump-file", default=None,
|
ap.add_argument("--dump-file", default=None,
|
||||||
help="путь к заранее скачанному дампу (надёжнее, чем читать из сети)")
|
help="путь к заранее скачанному дампу (надёжнее, чем читать из сети)")
|
||||||
|
ap.add_argument("--state-file", default="/tmp/bulk_wiki_state.json",
|
||||||
|
help="файл с числом уже обработанных статей — для продолжения после сбоя")
|
||||||
args = ap.parse_args()
|
args = ap.parse_args()
|
||||||
|
|
||||||
import psycopg2
|
import psycopg2
|
||||||
from app.algorithms.winnowing import winnow
|
from app.algorithms.winnowing import sample_evenly, winnow_ordered
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
|
||||||
conn = psycopg2.connect(host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
|
def connect():
|
||||||
|
return psycopg2.connect(
|
||||||
|
host=settings.POSTGRES_HOST, port=settings.POSTGRES_PORT,
|
||||||
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
|
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
|
||||||
password=settings.POSTGRES_PASSWORD)
|
password=settings.POSTGRES_PASSWORD,
|
||||||
|
)
|
||||||
|
|
||||||
|
conn = connect()
|
||||||
added = skipped = 0
|
added = skipped = 0
|
||||||
fp_total = 0
|
fp_total = 0
|
||||||
|
seen = 0 # сколько статей прошло через парсер (позиция в дампе)
|
||||||
t0 = time.time()
|
t0 = time.time()
|
||||||
batch: list[tuple[str, str, str]] = []
|
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
|
nonlocal added, skipped, fp_total
|
||||||
if not batch:
|
|
||||||
return
|
|
||||||
with conn.cursor() as cur:
|
with conn.cursor() as cur:
|
||||||
cur.executemany(
|
cur.executemany(
|
||||||
"""INSERT INTO documents (ext_id, source, title, lang, abstract, url,
|
"""INSERT INTO documents (ext_id, source, title, lang, abstract, url,
|
||||||
@@ -184,7 +203,9 @@ def main() -> None:
|
|||||||
if cur.fetchone():
|
if cur.fetchone():
|
||||||
skipped += 1
|
skipped += 1
|
||||||
continue
|
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):
|
for pos, h in enumerate(hashes):
|
||||||
buf.write(f"{doc_id}\t{h}\t{pos}\n")
|
buf.write(f"{doc_id}\t{h}\t{pos}\n")
|
||||||
fp_total += len(hashes)
|
fp_total += len(hashes)
|
||||||
@@ -195,11 +216,42 @@ def main() -> None:
|
|||||||
added += fresh
|
added += fresh
|
||||||
conn.commit()
|
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):
|
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))
|
batch.append((pid, title, text))
|
||||||
if len(batch) >= args.batch:
|
if len(batch) >= args.batch:
|
||||||
flush(batch)
|
flush(batch)
|
||||||
batch = []
|
batch = []
|
||||||
|
with open(args.state_file, "w") as fh:
|
||||||
|
json.dump({"processed": seen}, fh)
|
||||||
el = time.time() - t0
|
el = time.time() - t0
|
||||||
print(f" залито {added} · пропущено {skipped} · отпечатков {fp_total:,} · "
|
print(f" залито {added} · пропущено {skipped} · отпечатков {fp_total:,} · "
|
||||||
f"{added / max(el, 1):.1f} ст/с", flush=True)
|
f"{added / max(el, 1):.1f} ст/с", flush=True)
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ import io
|
|||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
import uuid
|
import uuid
|
||||||
from datetime import datetime, timedelta
|
from datetime import UTC, datetime, timedelta
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
@@ -925,6 +925,25 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
|
|||||||
)
|
)
|
||||||
).scalars().all()
|
).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 {
|
return {
|
||||||
"generated_at": datetime.utcnow().isoformat(timespec="seconds"),
|
"generated_at": datetime.utcnow().isoformat(timespec="seconds"),
|
||||||
"celery": celery_state,
|
"celery": celery_state,
|
||||||
@@ -956,6 +975,17 @@ async def debug_snapshot(db: AsyncSession = Depends(get_db)) -> dict:
|
|||||||
}
|
}
|
||||||
for t in failed_tasks
|
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) и влияют на
|
# Ключи бэкендов задаются в .env (генерируется из Infisical) и влияют на
|
||||||
# поведение воркеров — при разборе «почему не так считает» нужны первыми
|
# поведение воркеров — при разборе «почему не так считает» нужны первыми
|
||||||
"config": {
|
"config": {
|
||||||
|
|||||||
@@ -85,6 +85,68 @@ def winnow(text: str, k: int = 5, window: int = 4) -> set[int]:
|
|||||||
return fingerprint
|
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:
|
def jaccard_similarity(fp_a: set[int], fp_b: set[int]) -> float:
|
||||||
"""
|
"""
|
||||||
Коэффициент Жаккара для двух fingerprint'ов.
|
Коэффициент Жаккара для двух fingerprint'ов.
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ from sqlalchemy import func, select, update
|
|||||||
from sqlalchemy.exc import IntegrityError
|
from sqlalchemy.exc import IntegrityError
|
||||||
|
|
||||||
from app.algorithms.minhash import add_to_lsh, find_similar
|
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.celery_app import celery_app
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status
|
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
|
# Вычислить fingerprints
|
||||||
text = doc_data.get("full_text") or doc_data.get("abstract", "") or ""
|
text = doc_data.get("full_text") or doc_data.get("abstract", "") or ""
|
||||||
if text:
|
if text:
|
||||||
fp = winnow(text)
|
fingerprints_to_add = sample_evenly(
|
||||||
fingerprints_to_add = list(fp)[: settings.MAX_FINGERPRINTS_PER_DOC]
|
winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC
|
||||||
|
)
|
||||||
for i, hash_val in enumerate(fingerprints_to_add):
|
for i, hash_val in enumerate(fingerprints_to_add):
|
||||||
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
|
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",
|
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:
|
with db_session() as session:
|
||||||
doc = session.get(Document, doc_id)
|
doc = session.get(Document, doc_id)
|
||||||
|
|||||||
@@ -5,7 +5,9 @@ from app.algorithms.winnowing import (
|
|||||||
get_ngrams,
|
get_ngrams,
|
||||||
hash_ngram,
|
hash_ngram,
|
||||||
jaccard_similarity,
|
jaccard_similarity,
|
||||||
|
sample_evenly,
|
||||||
winnow,
|
winnow,
|
||||||
|
winnow_ordered,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Достаточно длинный текст, чтобы окно Winnowing реально отработало
|
# Достаточно длинный текст, чтобы окно Winnowing реально отработало
|
||||||
@@ -77,3 +79,57 @@ def test_compute_similarity_partial_overlap_is_between():
|
|||||||
modified = LONG + " добавлен ещё один совершенно новый хвост предложения здесь"
|
modified = LONG + " добавлен ещё один совершенно новый хвост предложения здесь"
|
||||||
sim = compute_similarity(LONG, modified)
|
sim = compute_similarity(LONG, modified)
|
||||||
assert 0.0 < sim < 1.0
|
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 # шаг ровный
|
||||||
|
|||||||
Reference in New Issue
Block a user