From 780a0a10613fda66ad6dfbe3b5a8b7ba1998f7a7 Mon Sep 17 00:00:00 2001 From: jze9 Date: Mon, 31 Aug 2026 20:54:59 +0500 Subject: [PATCH] =?UTF-8?q?fix(ops):=20=D1=87=D0=B8=D1=82=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=20=D0=B4=D0=B0=D0=BC=D0=BF=20=D0=92=D0=B8=D0=BA=D0=B8?= =?UTF-8?q?=D0=BF=D0=B5=D0=B4=D0=B8=D0=B8=20=D1=81=20=D0=B4=D0=B8=D1=81?= =?UTF-8?q?=D0=BA=D0=B0=20=E2=80=94=20=D1=81=D0=B5=D1=82=D0=B5=D0=B2=D0=BE?= =?UTF-8?q?=D0=B9=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=20=D1=80=D0=B2=D1=91?= =?UTF-8?q?=D1=82=D1=81=D1=8F=20=D0=BD=D0=B0=205.6=20=D0=93=D0=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wikimedia закрывает долгие потоковые соединения: обрыв пришёлся на 32 МБ из 5.9 ГБ. Скачать файл с докачкой (curl -C -) и читать локально надёжнее, поэтому у скрипта появился --dump-file; чтение из сети осталось запасным путём. Co-Authored-By: Claude Opus 5 --- scripts/ops/bulk_ingest_wikipedia_ru.py | 100 +++++++++++++++--------- 1 file changed, 65 insertions(+), 35 deletions(-) diff --git a/scripts/ops/bulk_ingest_wikipedia_ru.py b/scripts/ops/bulk_ingest_wikipedia_ru.py index d083546..49418fc 100755 --- a/scripts/ops/bulk_ingest_wikipedia_ru.py +++ b/scripts/ops/bulk_ingest_wikipedia_ru.py @@ -63,47 +63,75 @@ def clean_wikitext(raw: str) -> str: return text.strip() -def iter_pages(min_chars: int): - """Потоково читать дамп и отдавать статьи основного пространства имён.""" +def _pages_from_chunks(chunks, min_chars: int): + """Разобрать поток распакованного XML на статьи основного пространства имён.""" + buf = "" + for raw in chunks: + buf += raw.decode("utf-8", errors="replace") + + while True: + m = PAGE_RE.search(buf) + if not m: + break + page, buf = m.group(1), buf[m.end():] + + if REDIRECT_RE.search(page): + continue + ns = NS_RE.search(page) + if not ns or ns.group(1) != "0": # только статьи + continue + tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page) + if not (tm and im and xm): + continue + text = clean_wikitext(xm.group(1)) + if len(text) < min_chars: + continue + yield im.group(1), tm.group(1), text + + if len(buf) > 20 * 1024 * 1024: # страховка от разбухания + buf = buf[-1024 * 1024:] + + +def iter_pages(min_chars: int, dump_file: str | None): + """Статьи из локального дампа либо, если файла нет, прямо из сети. + + Локальный файл предпочтителен: Wikimedia рвёт долгие потоковые соединения + (проверено — обрыв на 32 МБ из 5.9 ГБ), а скачать файл можно с докачкой + (`curl -C -`) и потом читать сколько угодно. + """ + decomp = bz2.BZ2Decompressor() + + if dump_file: + def chunks(): + with open(dump_file, "rb") as fh: + while True: + part = fh.read(4 * 1024 * 1024) + if not part: + return + try: + yield decomp.decompress(part) + except EOFError: + return + yield from _pages_from_chunks(chunks(), min_chars) + return + import httpx - decomp = bz2.BZ2Decompressor() - buf = "" # Wikimedia отдаёт 403 без осмысленного User-Agent — по их правилам он должен # называть приложение и давать контакт headers = {"User-Agent": "AcademicHelper/1.0 (https://academic.jze9.ru; noreply@jze9.ru)"} - with httpx.stream("GET", DUMP_URL, timeout=120, follow_redirects=True, headers=headers) as resp: - resp.raise_for_status() - for chunk in resp.iter_bytes(4 * 1024 * 1024): - try: - raw = decomp.decompress(chunk) - except EOFError: - break - if not raw: - continue - buf += raw.decode("utf-8", errors="replace") - while True: - m = PAGE_RE.search(buf) - if not m: - break - page, buf = m.group(1), buf[m.end():] + def net_chunks(): + with httpx.stream("GET", DUMP_URL, timeout=120, follow_redirects=True, + headers=headers) as resp: + resp.raise_for_status() + for chunk in resp.iter_bytes(4 * 1024 * 1024): + try: + yield decomp.decompress(chunk) + except EOFError: + return - if REDIRECT_RE.search(page): - continue - ns = NS_RE.search(page) - if not ns or ns.group(1) != "0": # только статьи - continue - tm, im, xm = TITLE_RE.search(page), ID_RE.search(page), TEXT_RE.search(page) - if not (tm and im and xm): - continue - text = clean_wikitext(xm.group(1)) - if len(text) < min_chars: - continue - yield im.group(1), tm.group(1), text - - if len(buf) > 20 * 1024 * 1024: # страховка от разбухания - buf = buf[-1024 * 1024:] + yield from _pages_from_chunks(net_chunks(), min_chars) def main() -> None: @@ -113,6 +141,8 @@ def main() -> None: ap.add_argument("--batch", type=int, default=500, 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="путь к заранее скачанному дампу (надёжнее, чем читать из сети)") args = ap.parse_args() import psycopg2 @@ -165,7 +195,7 @@ def main() -> None: added += fresh conn.commit() - for pid, title, text in iter_pages(args.min_chars): + for pid, title, text in iter_pages(args.min_chars, args.dump_file): batch.append((pid, title, text)) if len(batch) >= args.batch: flush(batch)