diff --git a/scripts/ops/wikipedia_resume_offset.py b/scripts/ops/wikipedia_resume_offset.py new file mode 100755 index 0000000..7b501de --- /dev/null +++ b/scripts/ops/wikipedia_resume_offset.py @@ -0,0 +1,79 @@ +#!/usr/bin/env python3 +"""Вычислить байтовое смещение в дампе Википедии по уже залитым статьям. + +Нужен один раз при переходе на multistream-дамп: позиция продолжения сменила +смысл — раньше это был номер статьи (и прогон перечитывал дамп с начала), +теперь байтовое смещение блока. + +Как считается: берём максимальный page_id среди залитых статей и находим в +индексе дампа блок, которому он принадлежит. Страницы в дампе идут по +возрастанию page_id, поэтому всё, что дальше этого блока, ещё не залито. + +Запуск (в контейнере worker-indexer): + docker compose -f docker-compose.prod.yml exec -T worker-indexer python - \\ + < scripts/ops/wikipedia_resume_offset.py # показать + ... python - --apply < scripts/ops/wikipedia_resume_offset.py # записать +""" + +import argparse +import bz2 +import os + +DEFAULT_INDEX = "/parsers/ruwiki-index.txt.bz2" + + +def main() -> None: + ap = argparse.ArgumentParser(description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("--apply", action="store_true", help="записать смещение в источник") + ap.add_argument("--index", default=DEFAULT_INDEX, help="путь к индексу дампа") + args = ap.parse_args() + + from app.db import db_session + from sqlalchemy import text + + if not os.path.exists(args.index): + raise SystemExit(f"индекс не найден: {args.index}") + + with db_session() as s: + row = s.execute(text(""" + SELECT max(split_part(ext_id, ':', 2)::bigint) + FROM documents WHERE source = 'wikipedia_ru' + """)).first() + max_pageid = int(row[0]) if row and row[0] else 0 + print(f"максимальный page_id среди залитых: {max_pageid}") + + if not max_pageid: + print("залитых статей нет — начинать с нуля") + return + + # Индекс отсортирован по смещению, page_id внутри растут: ищем блок, + # в котором лежит наш максимальный, и берём его смещение + offset = 0 + found_title = "" + with bz2.open(args.index, "rt", encoding="utf-8") as fh: + for line in fh: + parts = line.split(":", 2) + if len(parts) < 3: + continue + off, pid = int(parts[0]), int(parts[1]) + if pid > max_pageid: + break + offset, found_title = off, parts[2].strip() + + print(f"блок с этой статьёй начинается на байте {offset} ({found_title[:50]!r})") + + if not args.apply: + print("\nзапустите с --apply, чтобы записать смещение в источник") + return + + with db_session() as s: + s.execute(text(""" + UPDATE parse_sources SET resume_token = :off WHERE source_type = 'wikipedia_ru' + """), {"off": str(offset)}) + s.commit() + print(f"смещение {offset} записано в источник") + + +if __name__ == "__main__": + main() diff --git a/scripts/parsers/wikipedia_ru.py b/scripts/parsers/wikipedia_ru.py index 5bcea3a..c3998c9 100644 --- a/scripts/parsers/wikipedia_ru.py +++ b/scripts/parsers/wikipedia_ru.py @@ -8,14 +8,22 @@ Дамп — 5.6 ГБ и около 2 млн статей, держать их в памяти нельзя, поэтому `run_parser` читает результат лениво и пишет пачками (см. bulk_writer.py). -Дамп качается заранее и кладётся туда, где его видит воркер: - curl -L -C - --retry 100 -A "AcademicHelper/1.0 (noreply@jze9.ru)" \\ - -o scripts/parsers/ruwiki.xml.bz2 \\ - https://dumps.wikimedia.org/ruwiki/latest/ruwiki-latest-pages-articles.xml.bz2 +Используется **multistream**-вариант дампа: он состоит из независимых bz2-блоков +по 100 статей, и к нему прилагается индекс со смещениями. Это принципиально — +обычный дамп читается только с начала, поэтому каждый следующий прогон +перечитывал всё уже залитое: на 20 тысячах статей это стоило 6 минут из 25 +доступных, а на 100 тысячах съело бы весь бюджет и заливка встала бы совсем. +С multistream позиция продолжения — байтовое смещение, и прогон стартует +мгновенно. + +Файлы качаются заранее и кладутся туда, где их видит воркер: + B=https://dumps.wikimedia.org/ruwiki/latest + UA="AcademicHelper/1.0 (https://academic.jze9.ru; noreply@jze9.ru)" + curl -L -C - --retry 100 -A "$UA" -o scripts/parsers/ruwiki-multistream.xml.bz2 \\ + $B/ruwiki-latest-pages-articles-multistream.xml.bz2 Читать дамп прямо из сети не выйдет: Wikimedia обрывает долгие соединения -(проверено — обрыв на 32 МБ из 5.9 ГБ), да и продолжить с места разрыва -распакованный поток bz2 нельзя. +(проверено — обрыв на 32 МБ из 5.9 ГБ). """ import bz2 @@ -29,7 +37,7 @@ from base import BaseParser, ProgressCallback logger = logging.getLogger(__name__) -DEFAULT_DUMP = "/parsers/ruwiki.xml.bz2" +DEFAULT_DUMP = "/parsers/ruwiki-multistream.xml.bz2" PAGE_RE = re.compile(r"(.*?)", re.DOTALL) TITLE_RE = re.compile(r"(.*?)", re.DOTALL) @@ -78,18 +86,19 @@ class WikipediaRuParser(BaseParser): limit: int = 1000, min_chars: int = 2000, dump_path: str | None = None, - skip: int = 0, + start_offset: int = 0, progress_cb: ProgressCallback | None = None, **_ignored: Any, ) -> Iterator[dict[str, Any]]: - """Статьи основного пространства имён из дампа. + """Статьи основного пространства имён из multistream-дампа. Args: limit: сколько статей отдать min_chars: минимальная длина текста — заготовки в корпусе бесполезны - dump_path: путь к дампу (по умолчанию /parsers/ruwiki.xml.bz2) - skip: сколько статей пропустить — так прогон продолжается с места - обрыва, не перечитывая уже залитое + dump_path: путь к дампу (по умолчанию /parsers/ruwiki-multistream.xml.bz2) + start_offset: байтовое смещение в файле, с которого продолжать. + Именно смещение, а не номер статьи: перечитывание дампа с начала + росло линейно и на сотне тысяч статей съедало весь бюджет прогона progress_cb: см. base.ProgressCallback """ path = dump_path or os.environ.get("WIKIPEDIA_DUMP_PATH") or DEFAULT_DUMP @@ -99,23 +108,35 @@ class WikipediaRuParser(BaseParser): f"или укажите WIKIPEDIA_DUMP_PATH" ) - decomp = bz2.BZ2Decompressor() + given = 0 buf = "" - seen = 0 # статей прошло через парсер — позиция в дампе - given = 0 # статей отдано наружу with open(path, "rb") as fh: + fh.seek(start_offset) + # Смещение блока, из которого пришли уже отданные статьи: его и + # сохраняем как позицию продолжения, чтобы ничего не потерять + block_offset = start_offset + decomp = bz2.BZ2Decompressor() + while given < limit: part = fh.read(4 * 1024 * 1024) if not part: break - try: - raw = decomp.decompress(part) - except EOFError: - break - if not raw: - continue - buf += raw.decode("utf-8", errors="replace") + + # В multistream-дампе потоки идут подряд: закончился один — + # начинаем следующий с того места, где предыдущий остановился + while part: + try: + raw = decomp.decompress(part) + except (OSError, EOFError): + return + if raw: + buf += raw.decode("utf-8", errors="replace") + + if not decomp.eof: + break + part = decomp.unused_data + decomp = bz2.BZ2Decompressor() while given < limit: m = PAGE_RE.search(buf) @@ -127,18 +148,17 @@ class WikipediaRuParser(BaseParser): if doc is None: continue - seen += 1 - if seen <= skip: # быстро проматываем уже обработанное - continue - given += 1 - doc["dump_position"] = seen + doc["dump_offset"] = block_offset yield doc if progress_cb and not progress_cb(given): logger.info("Википедия: выборка остановлена по запросу (%d)", given) return + # Всё разобранное отдано — следующая позиция продолжения здесь + block_offset = fh.tell() - len(buf.encode("utf-8", errors="ignore")) // 4 + if len(buf) > 20 * 1024 * 1024: # страховка от разбухания buf = buf[-1024 * 1024:] @@ -175,5 +195,5 @@ class WikipediaRuParser(BaseParser): "url": f"https://ru.wikipedia.org/?curid={pid}", "abstract": (raw.get("text") or "")[:2000], "text": raw.get("text", ""), - "dump_position": raw.get("dump_position"), + "dump_offset": raw.get("dump_offset"), } diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index 094e1c0..da96cba 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -580,12 +580,14 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A "year_to": cfg.get("year_to"), } elif source_type == "wikipedia_ru": - # Массовый источник: дамп читается с позиции прошлого прогона + # Массовый источник: позиция — байтовое смещение в multistream-дампе. + # Раньше хранился номер статьи, и каждый прогон перечитывал дамп с + # начала: на 20 тысячах это стоило 6 минут из 25, дальше росло линейно from wikipedia_ru import WikipediaRuParser as P kwargs = { "limit": limit, "dump_path": query or None, # query = путь к дампу, если задан - "skip": int(cfg.get("resume_token") or 0), + "start_offset": int(cfg.get("resume_token") or 0), } elif source_type == "pmc_bulk": from pmc_bulk import PMCBulkParser as P @@ -742,7 +744,7 @@ def _ingest_bulk( # Позиция продолжения: у Википедии — номер статьи в дампе, # у PMC — токен страницы бакета - resume_token = doc.get("dump_position") or doc.get("resume_token") or resume_token + resume_token = doc.get("dump_offset") or doc.get("resume_token") or resume_token batch.append(doc) prog.count_fetched(prog.fetched + 1)