"""Большой бэг `new_data` кусками прямо из tar-архива. Это 90 ГБ в 221 шарде, а свободного места на диске меньше, чем весь архив. Кусок из нескольких подряд идущих шардов распаковывается во временный каталог, читается как обычный многошардовый бэг (`flyguard.bag.Bag`) и удаляется. Шард достаётся по смещению в архиве, без повторного разбора заголовков, поэтому куски можно распаковывать из нескольких процессов сразу. """ from __future__ import annotations import os import re import shutil import tarfile import tempfile from contextlib import contextmanager from pathlib import Path # Архив лежит там, куда его положили при скачивании датасета, а не в data/: # распаковывать его целиком некуда. Путь переопределяется FLYGUARD_NEW_DATA # или ключом --tar у инструментов. DEFAULT_TAR = os.environ.get( "FLYGUARD_NEW_DATA", str(Path(__file__).resolve().parents[1] / "датасет" / "new_data")) SHARD_RE = re.compile(r"_(\d+)\.db3$") def shards(tar_path: str | Path) -> list[tuple[int, str, int, int]]: """(номер, имя файла, смещение данных, размер) всех шардов, по номеру.""" out = [] with tarfile.open(tar_path, "r:") as t: for m in t: mm = SHARD_RE.search(m.name) if mm and m.isfile(): out.append((int(mm.group(1)), Path(m.name).name, m.offset_data, m.size)) out.sort() return out def pick(members: list, spec: str) -> list: """Шарды по срезу номеров: `110:` — со 110-го до конца, `0:110` — первые 110.""" lo, _, hi = spec.partition(":") lo_i = int(lo) if lo else 0 hi_i = int(hi) if hi else None return [m for m in members if m[0] >= lo_i and (hi_i is None or m[0] < hi_i)] def split(members: list, per_chunk: int) -> list[list]: """Подряд идущие шарды группами: каждая группа — одна «запись».""" return [members[i:i + per_chunk] for i in range(0, len(members), per_chunk)] @contextmanager def chunk(tar_path: str | Path, members: list, workdir: str | None = None): """Распаковать шарды во временный каталог, отдать его путь, потом удалить.""" d = Path(tempfile.mkdtemp(prefix="fg_nd_", dir=workdir)) try: with open(tar_path, "rb") as src: for _, name, off, size in members: src.seek(off) with open(d / name, "wb") as dst: left = size while left: buf = src.read(min(left, 1 << 22)) if not buf: raise OSError(f"архив обрезан на {name}") dst.write(buf) left -= len(buf) yield d finally: shutil.rmtree(d, ignore_errors=True)