From f4092e53312bcd19248dfe57383bc90754052992 Mon Sep 17 00:00:00 2001 From: jze9 Date: Sat, 5 Sep 2026 19:13:01 +0500 Subject: [PATCH] =?UTF-8?q?fix(ops):=20=D0=B7=D0=B0=D0=BB=D0=B8=D0=B2?= =?UTF-8?q?=D0=BA=D0=B8=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B6=D0=B8=D0=B2=D0=B0?= =?UTF-8?q?=D1=8E=D1=82=20=D0=BE=D0=B1=D1=80=D1=8B=D0=B2=D1=8B=20=D0=B1?= =?UTF-8?q?=D0=B0=D0=B7=D1=8B=20=D0=B8=20=D1=81=D0=B5=D1=82=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit За неделю обе массовые заливки умерли молча и не возобновились: Википедия на 20 008 статьях из ~2 млн («server closed the connection»), PMC примерно на 85 тыс. из 100 тыс. (таймаут SSL-рукопожатия). Ни ретраев, ни возобновления. - запись пачки повторяется с переподключением к базе (5 попыток с паузой); - обрыв листинга бакета стоит паузы, а не всей заливки; - у Википедии сохраняется позиция в дампе — после сбоя продолжаем с неё, а не проматываем с нуля, полагаясь на дедуп по ext_id; - батч уменьшен (500 статей × 2000 отпечатков = миллион строк в одной транзакции — вероятная причина обрыва); - обрезка отпечатков переведена на равномерную выборку. Co-Authored-By: Claude Opus 5 --- scripts/ops/bulk_ingest_pmc.py | 158 +++++++++++++++--------- scripts/ops/bulk_ingest_wikipedia_ru.py | 70 +++++++++-- 2 files changed, 163 insertions(+), 65 deletions(-) diff --git a/scripts/ops/bulk_ingest_pmc.py b/scripts/ops/bulk_ingest_pmc.py index 0bff506..042b229 100755 --- a/scripts/ops/bulk_ingest_pmc.py +++ b/scripts/ops/bulk_ingest_pmc.py @@ -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__": diff --git a/scripts/ops/bulk_ingest_wikipedia_ru.py b/scripts/ops/bulk_ingest_wikipedia_ru.py index 49418fc..7625703 100755 --- a/scripts/ops/bulk_ingest_wikipedia_ru.py +++ b/scripts/ops/bulk_ingest_wikipedia_ru.py @@ -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)