forked from Dan4ick/Lidar_Muxa
72 lines
3.3 KiB
Python
72 lines
3.3 KiB
Python
"""Большой бэг `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)
|