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>
This commit is contained in:
@@ -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)
|
||||||
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:
|
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():
|
||||||
dbname=settings.POSTGRES_DB, user=settings.POSTGRES_USER,
|
return psycopg2.connect(
|
||||||
password=settings.POSTGRES_PASSWORD)
|
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
|
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)
|
||||||
|
|||||||
Reference in New Issue
Block a user