feat(ingestion): скрипт бэкфилла full_text для CyberLeninka из OCR

Оставался незакоммиченным с предыдущего шага (докачка full-text). Переспрашивает
59 уже засеянных cyberleninka-тем из parse_sources, матчит raw-статьи с уже
залитыми documents по ext_id, и там где minio_key ещё не проставлен — пересчитывает
Winnowing-отпечатки из OCR full_text (не короткой annotation) и сохраняет текст в
MinIO. Идемпотентно (пропускает уже обработанные). Не трогает faiss_id/эмбеддинги.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
jze9
2026-08-24 15:09:18 +05:00
parent 2aac40ed3e
commit 92a1ce0f57

View File

@@ -0,0 +1,143 @@
#!/usr/bin/env python3
"""Бэкфилл full_text для уже залитых документов CyberLeninka из OCR-фрагментов.
Парсер (scripts/parsers/cyberleninka.py) раньше игнорировал поле "ocr" в ответе
поиска (список OCR-фрагментов текста статьи) — full_text всегда был None, и
Winnowing-отпечатки (L1) считались по короткой аннотации. Это исправлено для
НОВЫХ документов; этот скрипт досчитывает уже существующие.
Алгоритм: для каждой ранее засеянной cyberleninka-темы (parse_sources) —
переспросить search API (тот же query/limit — переспрос идемпотентен и не бьёт
по документу дважды), сматчить raw-статьи с уже существующими documents по
ext_id, и там где full_text ещё не сохранён (minio_key IS NULL):
1. пересчитать fingerprints по full_text (заменить провизорные из annotation);
2. сохранить full_text в MinIO (тот же путь, что enrich_full_text: corpus/{id}.txt);
3. проставить documents.minio_key.
Не трогает faiss_id/эмбеддинги (embed_documents использует title+abstract, не
full_text — пересчитывать нечего). Идемпотентно: документы с уже проставленным
minio_key пропускаются.
Запуск:
python scripts/backfill_cyberleninka_fulltext.py # dry-run (посчитать)
python scripts/backfill_cyberleninka_fulltext.py --apply # применить
"""
import argparse
import io
import os
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent / "parsers"))
sys.path.insert(0, str(Path(__file__).resolve().parent.parent
/ "services" / "worker-indexer" / "app" / "algorithms"))
def _load_env() -> None:
env = Path(__file__).resolve().parent.parent / ".env"
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
k, _, v = line.partition("=")
os.environ.setdefault(k.strip(), v.strip())
def _pg_dsn() -> dict:
return {
"host": os.environ["POSTGRES_HOST"], "port": int(os.environ.get("POSTGRES_PORT", "5432")),
"dbname": os.environ["POSTGRES_DB"], "user": os.environ["POSTGRES_USER"],
"password": os.environ["POSTGRES_PASSWORD"],
}
MAX_FP = 20000
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--apply", action="store_true")
args = ap.parse_args()
_load_env()
import psycopg2
from cyberleninka import CyberLeninkaParser
from minio import Minio
from winnowing import winnow
conn = psycopg2.connect(**_pg_dsn())
minio = Minio(
os.environ["MINIO_ENDPOINT"],
access_key=os.environ["MINIO_ACCESS_KEY"], secret_key=os.environ["MINIO_SECRET_KEY"],
secure=False,
)
bucket = os.environ.get("MINIO_BUCKET_DOCS", "documents")
with conn.cursor() as cur:
cur.execute(
"SELECT id, query, \"limit\" FROM parse_sources "
"WHERE source_type='cyberleninka' AND last_status='done' ORDER BY id"
)
topics = cur.fetchall()
print(f"Тем CyberLeninka к обработке: {len(topics)}")
parser = CyberLeninkaParser()
matched = updated = no_ocr = not_found = already_done = 0
for _sid, query, limit in topics:
raws = list(parser.fetch(query=query or "", limit=limit))
for raw in raws:
d = parser.transform(raw)
if not (d and d.get("ext_id")):
continue
with conn.cursor() as cur:
cur.execute(
"SELECT id, minio_key FROM documents WHERE ext_id=%s AND source='cyberleninka'",
(d["ext_id"],),
)
row = cur.fetchone()
if not row:
not_found += 1
continue
doc_id, minio_key = row
matched += 1
if minio_key:
already_done += 1
continue
if not d.get("full_text"):
no_ocr += 1
continue
if args.apply:
text = d["full_text"]
key = f"corpus/{doc_id}.txt"
data = text.encode("utf-8")
minio.put_object(bucket, key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8")
hashes = list(winnow(text))[:MAX_FP]
with conn.cursor() as cur:
cur.execute("DELETE FROM fingerprints WHERE doc_id=%s", (doc_id,))
if hashes:
from psycopg2.extras import execute_values
execute_values(
cur, "INSERT INTO fingerprints (doc_id,hash_value,position) VALUES %s",
[(doc_id, h, i) for i, h in enumerate(hashes)],
)
cur.execute("UPDATE documents SET minio_key=%s WHERE id=%s", (key, doc_id))
conn.commit()
updated += 1
print(f" {(query or '')[:32]:32} raw={len(raws):4} | обновлено всего={updated}", flush=True)
print(f"\nИТОГ: сматчено={matched} обновлено={updated} "
f"уже_было={already_done} без_ocr={no_ocr} не_найдено={not_found}"
f"{' [DRY-RUN — ничего не записано, добавьте --apply]' if not args.apply else ''}")
conn.close()
if __name__ == "__main__":
main()