параллельный прогон инструментов по бэгам
This commit is contained in:
parent
d76c99eb60
commit
00e61bcabc
6 changed files with 869 additions and 660 deletions
13
README.md
13
README.md
|
|
@ -87,7 +87,7 @@ flyguard/ ядро: стадии обработки, память, сч
|
|||
pipeline.py сборка
|
||||
synth.py вставка предметов трассировкой лучей
|
||||
tools/ обучение, оценка, разбор
|
||||
tests/ 27 тестов, запускаются без данных и без ROS
|
||||
tests/ 29 тестов, запускаются без данных и без ROS
|
||||
docs/ методика и результаты
|
||||
artifacts/ обученные модели
|
||||
```
|
||||
|
|
@ -119,6 +119,17 @@ python tools/make_benchmark.py --memory artifacts/mushroom_body.npz \
|
|||
python tools/plot_benchmark.py # кривые и график
|
||||
```
|
||||
|
||||
Тяжёлые шаги сами раскладываются по бэгам на процессы — записей пять, физических
|
||||
ядер шесть, и это вся доступная зернистость: конвейер держит состояние между
|
||||
кадрами, поэтому разрезать одну запись нельзя. Замерено: полигон 134 → 36 с,
|
||||
сбор выборки 96 → 26 с на облегчённой конфигурации, то есть 3.7–3.8×, и файл на
|
||||
выходе совпадает с последовательным **побайтово**. Отключается `--jobs 1`.
|
||||
|
||||
Для замера задержки кадра `--jobs 1` обязателен: под пятью процессами время
|
||||
кадра растёт с 32 до 56 мс. Это свойство замера, а не конвейера, поэтому
|
||||
`evaluate.py` в параллельном режиме печатает задержку как `nan` — чтобы такое
|
||||
число нельзя было случайно привести в отчёте.
|
||||
|
||||
---
|
||||
|
||||
## Где мы сейчас
|
||||
|
|
|
|||
|
|
@ -404,3 +404,49 @@ def test_real_bag_projects_without_angular_error():
|
|||
r = img.r_near[img.valid]
|
||||
assert r.min() > 0 and r.max() < 250
|
||||
assert np.all(img.r_far[img.valid] >= img.r_near[img.valid] - 1e-3)
|
||||
|
||||
|
||||
# ------------------------------------------------------------ раскладка по бэгам
|
||||
|
||||
_TOOLS = ROOT / "tools" # выгрузка: tools лежит рядом с тестами
|
||||
if not _TOOLS.exists():
|
||||
_TOOLS = ROOT.parents[2] / "tools" # основной проект: ros2_ws/src/flyguard
|
||||
if str(_TOOLS) not in sys.path:
|
||||
sys.path.insert(0, str(_TOOLS))
|
||||
|
||||
import _parallel as _P # noqa: E402
|
||||
|
||||
|
||||
def _twice(x):
|
||||
"""Задача для проверки. Верхнего уровня: иначе её не передать в процесс."""
|
||||
return x * 2
|
||||
|
||||
|
||||
def test_parallel_keeps_task_order_when_results_arrive_out_of_order():
|
||||
"""Считается по готовности, складывается по номеру задачи.
|
||||
|
||||
Ломается это незаметно и опасно: цифры остаются правдоподобными, просто
|
||||
приписываются не тому бэгу. Поэтому проверяется не «столько же строк», а
|
||||
что результат каждой задачи лёг на своё место.
|
||||
"""
|
||||
tasks = list(range(7))
|
||||
want = [(t, t * 2) for t in tasks]
|
||||
|
||||
seq = [None] * len(tasks)
|
||||
for i, t, r, _ in _P.run(_twice, tasks, jobs=1):
|
||||
seq[i] = (t, r)
|
||||
|
||||
par = [None] * len(tasks)
|
||||
for i, t, r, _ in _P.run(_twice, tasks, jobs=4):
|
||||
par[i] = (t, r)
|
||||
|
||||
assert seq == want
|
||||
assert par == want
|
||||
|
||||
|
||||
def test_parallel_never_starts_more_processes_than_there_are_bags():
|
||||
"""Бэгов пять, и шестой процесс занять нечем."""
|
||||
assert _P.resolve(0, 1) == 1 # одна задача — без пула вовсе
|
||||
assert _P.resolve(8, 3) == 3 # просили больше, чем есть работы
|
||||
assert _P.resolve(1, 5) == 1 # явная последовательная отладка
|
||||
assert 1 <= _P.resolve(0, 5) <= 5
|
||||
|
|
|
|||
91
tools/_parallel.py
Normal file
91
tools/_parallel.py
Normal file
|
|
@ -0,0 +1,91 @@
|
|||
"""Раскладка задач по бэгам на процессы.
|
||||
|
||||
Конвейер держит состояние между кадрами — пройденный путь, треки, накопитель
|
||||
складчатого тела, — поэтому разрезать один бэг нельзя. Зато бэги независимы
|
||||
друг от друга, и это ровно та зернистость, которая нужна: пять записей на
|
||||
шесть физических ядер.
|
||||
|
||||
Результат совпадает с последовательным прогоном **точно**, а не «примерно»:
|
||||
генератор случайных чисел создаётся внутри задачи от того же зерна, общей
|
||||
изменяемой памяти между бэгами нет. Порядок родитель восстанавливает сам —
|
||||
печатает по готовности, а сохраняет в исходном.
|
||||
|
||||
Ускорение упирается не в ядра, а в число бэгов: их пять, и больше пяти
|
||||
процессов тут просто нечем занять.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import time
|
||||
from concurrent.futures import ProcessPoolExecutor, as_completed
|
||||
|
||||
|
||||
def limit_threads() -> None:
|
||||
"""Один поток BLAS на процесс.
|
||||
|
||||
Вызывать ДО создания пула: переменные наследуются дочерним процессом при
|
||||
запуске, а число потоков BLAS выбирает один раз при импорте numpy и потом
|
||||
не меняет. Без этого каждый из пяти воркеров разворачивается на все ядра
|
||||
и они дерутся за те же шесть.
|
||||
"""
|
||||
for v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS",
|
||||
"NUMEXPR_NUM_THREADS", "VECLIB_MAXIMUM_THREADS"):
|
||||
os.environ.setdefault(v, "1")
|
||||
|
||||
|
||||
def resolve(jobs: int, n_tasks: int) -> int:
|
||||
"""Сколько процессов поднимать: 0 — по числу задач, но не больше ядер."""
|
||||
if n_tasks <= 1:
|
||||
return 1
|
||||
if jobs > 0:
|
||||
return max(1, min(jobs, n_tasks))
|
||||
cpu = os.cpu_count() or 2
|
||||
return max(1, min(n_tasks, cpu // 2)) # логических ядер вдвое больше физических
|
||||
|
||||
|
||||
class _Timed:
|
||||
"""Часы держит сам воркер.
|
||||
|
||||
В родителе видно только момент готовности, а он у всех задач, стартовавших
|
||||
разом, почти один и тот же — по такому замеру не понять, какой бэг тяжёлый.
|
||||
"""
|
||||
|
||||
def __init__(self, fn):
|
||||
self.fn = fn
|
||||
|
||||
def __call__(self, task):
|
||||
t0 = time.time()
|
||||
return self.fn(task), time.time() - t0
|
||||
|
||||
|
||||
def run(work, tasks, jobs: int):
|
||||
"""Выполнить `work(task)` по всем задачам, отдавая `(i, task, res, с)`.
|
||||
|
||||
Отдаёт по мере готовности, поэтому `i` — исходный номер задачи, и по нему
|
||||
вызывающий раскладывает результаты обратно в порядок бэгов.
|
||||
|
||||
При одном процессе всё считается прямо здесь, без пула: остаётся чем
|
||||
отлаживать, и трассировка ошибки не проходит через межпроцессную передачу.
|
||||
"""
|
||||
tasks = list(tasks)
|
||||
n = resolve(jobs, len(tasks))
|
||||
if n <= 1:
|
||||
for i, t in enumerate(tasks):
|
||||
t0 = time.time()
|
||||
yield i, t, work(t), time.time() - t0
|
||||
return
|
||||
|
||||
limit_threads()
|
||||
with ProcessPoolExecutor(max_workers=n) as ex:
|
||||
fut = {ex.submit(_Timed(work), t): (i, t) for i, t in enumerate(tasks)}
|
||||
for f in as_completed(fut):
|
||||
i, t = fut[f]
|
||||
res, secs = f.result()
|
||||
yield i, t, res, secs
|
||||
|
||||
|
||||
def add_argument(ap) -> None:
|
||||
"""Один и тот же флаг во всех инструментах, чтобы не помнить разные."""
|
||||
ap.add_argument("--jobs", type=int, default=0,
|
||||
help="сколько бэгов считать разом; 0 — по числу бэгов, "
|
||||
"но не больше физических ядер; 1 — в один процесс")
|
||||
|
|
@ -15,11 +15,12 @@ from __future__ import annotations
|
|||
|
||||
import argparse
|
||||
import json
|
||||
import time
|
||||
import tempfile
|
||||
|
||||
import numpy as np
|
||||
|
||||
import _bootstrap as B # noqa: F401
|
||||
import _parallel as P
|
||||
from flyguard.bag import Bag, find_bags
|
||||
from flyguard.mushroom_body import MushroomBody, MushroomBodyConfig
|
||||
from flyguard.pipeline import FlyGuard, Params
|
||||
|
|
@ -28,6 +29,17 @@ OBSTACLE_BAG = "doubleT_obstacle"
|
|||
TRUE_D = (50.0, 62.0)
|
||||
|
||||
|
||||
def _work(task):
|
||||
"""Одна задача — один бэг. Память и считывание грузятся по пути уже здесь."""
|
||||
path, params, limit, mem_path, rd_path = task
|
||||
memory = MushroomBody.load(mem_path) if mem_path else None
|
||||
readout = None
|
||||
if rd_path:
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
readout = MbonReadout.load(rd_path)
|
||||
return run_bag(path, memory, limit, params, readout=readout)
|
||||
|
||||
|
||||
def train_excluding(per_bag: dict[str, np.ndarray], extra: np.ndarray | None,
|
||||
exclude: str, target: float, device: str) -> MushroomBody:
|
||||
parts = [v for k, v in per_bag.items() if k not in (exclude, OBSTACLE_BAG)]
|
||||
|
|
@ -139,6 +151,7 @@ def main() -> None:
|
|||
ap.add_argument("--no-memory", action="store_true",
|
||||
help="совсем без памяти тоннеля — так выглядит первый проезд по новой линии")
|
||||
ap.add_argument("--out", default=str(B.ARTIFACTS / "generalisation.json"))
|
||||
P.add_argument(ap)
|
||||
args = ap.parse_args()
|
||||
|
||||
d = np.load(args.cache, allow_pickle=True)
|
||||
|
|
@ -198,35 +211,57 @@ def main() -> None:
|
|||
if args.mbon_prior_from is not None:
|
||||
over["mbon_prior_from"] = args.mbon_prior_from
|
||||
params = Params(**over)
|
||||
readout = None
|
||||
folds = {}
|
||||
fold_paths: dict[str, str] = {}
|
||||
if args.mbon_dir:
|
||||
from pathlib import Path as _P
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
for f in _P(args.mbon_dir).glob("mbon_*.npz"):
|
||||
folds[f.stem[len("mbon_"):]] = MbonReadout.load(f)
|
||||
fold_paths[f.stem[len("mbon_"):]] = str(f)
|
||||
print(f"считывание MBON по складкам: {args.mbon_dir} "
|
||||
f"({len(folds)} моделей)")
|
||||
f"({len(fold_paths)} моделей)")
|
||||
elif args.mbon:
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
readout = MbonReadout.load(args.mbon)
|
||||
print(f"считывание MBON: {args.mbon}")
|
||||
rows = []
|
||||
|
||||
# Память обучается ЗДЕСЬ, а не в воркере: обучение идёт на видеокарте, и
|
||||
# делить одну карту на пять процессов незачем. Воркеру достаётся готовый
|
||||
# файл — долгая часть это проход по записи, а не обучение.
|
||||
tmp = None if args.no_memory else tempfile.TemporaryDirectory(prefix="fg_mem_")
|
||||
tasks, sizes = [], []
|
||||
for p in find_bags(args.root):
|
||||
t0 = time.time()
|
||||
mem = (None if args.no_memory else
|
||||
train_excluding(per_bag, extra, p.name, args.target, args.device))
|
||||
rd = folds.get(p.name, readout) if folds else readout
|
||||
if folds and p.name not in folds and p.name != OBSTACLE_BAG:
|
||||
mem_path, n_seen = "", 0
|
||||
if tmp is not None:
|
||||
mem = train_excluding(per_bag, extra, p.name, args.target, args.device)
|
||||
mem_path = str(Path(tmp.name) / f"mem_{p.name}.npz")
|
||||
mem.save(mem_path)
|
||||
n_seen = int(mem.n_seen)
|
||||
del mem
|
||||
rd = fold_paths.get(p.name, args.mbon) if fold_paths else args.mbon
|
||||
if fold_paths and p.name not in fold_paths and p.name != OBSTACLE_BAG:
|
||||
print(f" внимание: для {p.name} нет своей складки")
|
||||
r = run_bag(p, mem, args.limit, params, readout=rd)
|
||||
r["train_size"] = int(mem.n_seen) if mem is not None else 0
|
||||
rows.append(r)
|
||||
tasks.append((p, params, args.limit, mem_path, rd))
|
||||
sizes.append(n_seen)
|
||||
|
||||
# Задержка кадра под пятью процессами вырастает с 32 до 56 мс — это
|
||||
# свойство замера, а не конвейера, и такое число нельзя показывать как
|
||||
# запас по бюджету. Поэтому в параллельном режиме оно не печатается
|
||||
# вовсе: перепутать nan с честным замером невозможно, а предупреждение
|
||||
# в шапке пролистывается.
|
||||
par = P.resolve(args.jobs, len(tasks)) > 1
|
||||
if par:
|
||||
print("параллельно: задержка кадра не измеряется, для неё нужен --jobs 1")
|
||||
slots: list = [None] * len(tasks)
|
||||
for i, task, r, secs in P.run(_work, tasks, args.jobs):
|
||||
r["train_size"] = sizes[i]
|
||||
if par:
|
||||
r["ms_p50"] = r["ms_p95"] = float("nan")
|
||||
slots[i] = r
|
||||
obj = f"объект {r['obj_rate']:6.1%} | " if r["obj_rate"] is not None else ""
|
||||
print(f"{r['bag']:40s} кадров {r['frames']:4d} путь {r['path_m']:6.0f} м | "
|
||||
f"{obj}тревог {r['alarm_rate']:6.1%} | ложных треков {r['fp_tracks']:3d} "
|
||||
f"({r['fp_per_km']:6.1f} на км) | {r['ms_p50']:5.1f}/{r['ms_p95']:5.1f} мс | "
|
||||
f"{time.time()-t0:5.0f} с", flush=True)
|
||||
f"{secs:5.0f} с", flush=True)
|
||||
rows = [r for r in slots if r is not None]
|
||||
if tmp is not None:
|
||||
tmp.cleanup()
|
||||
|
||||
B.ARTIFACTS.mkdir(parents=True, exist_ok=True)
|
||||
with open(args.out, "w", encoding="utf-8") as f:
|
||||
|
|
|
|||
|
|
@ -17,11 +17,11 @@ from __future__ import annotations
|
|||
|
||||
import argparse
|
||||
import json
|
||||
import time
|
||||
|
||||
import numpy as np
|
||||
|
||||
import _bootstrap as B # noqa: F401
|
||||
import _parallel as P
|
||||
from flyguard.bag import Bag, find_bags
|
||||
from flyguard.mushroom_body import MushroomBody
|
||||
from flyguard.pipeline import FlyGuard, Params
|
||||
|
|
@ -30,6 +30,22 @@ from flyguard.synth import IntensityEnv, Placement, catalogue, inject
|
|||
HOLDOUT = "doubleT_obstacle" # там уже есть настоящий объект
|
||||
|
||||
|
||||
def _work(task):
|
||||
"""Одна задача — один бэг.
|
||||
|
||||
Обученное грузится путями и уже внутри процесса: передавать модели через
|
||||
межпроцессную границу незачем, а свою складку каждый воркер берёт сам.
|
||||
"""
|
||||
path, params, limit, d_start, laterals, seed, mem_path, rd_path = task
|
||||
memory = MushroomBody.load(mem_path) if mem_path else None
|
||||
readout = None
|
||||
if rd_path:
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
readout = MbonReadout.load(rd_path)
|
||||
return run_bag(path, params, memory, limit, d_start, laterals, seed,
|
||||
readout=readout)
|
||||
|
||||
|
||||
def ego_track(bag: Bag, params: Params, limit: int | None):
|
||||
"""Первый проход: пройденный путь на каждом кадре (разметка по дистанции)."""
|
||||
fg = FlyGuard(params, memory=None)
|
||||
|
|
@ -155,22 +171,18 @@ def main() -> None:
|
|||
help="к какой дальности вклад знакомости обнуляется")
|
||||
ap.add_argument("--mbon-prior-from", type=float, default=None,
|
||||
help="с какой дальности поправлять оценку модели на распространённость предметов; 0 — не поправлять")
|
||||
P.add_argument(ap)
|
||||
args = ap.parse_args()
|
||||
|
||||
B.ARTIFACTS.mkdir(parents=True, exist_ok=True)
|
||||
memory = MushroomBody.load(args.memory) if args.memory else None
|
||||
readout = None
|
||||
folds: dict = {}
|
||||
fold_paths: dict[str, str] = {}
|
||||
if args.mbon_dir:
|
||||
from pathlib import Path as _P
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
for f in _P(args.mbon_dir).glob("mbon_*.npz"):
|
||||
folds[f.stem[len("mbon_"):]] = MbonReadout.load(f)
|
||||
fold_paths[f.stem[len("mbon_"):]] = str(f)
|
||||
print(f"считывание MBON по складкам: {args.mbon_dir} "
|
||||
f"({len(folds)} моделей)")
|
||||
f"({len(fold_paths)} моделей)")
|
||||
elif args.mbon:
|
||||
from flyguard.mbon_readout import MbonReadout
|
||||
readout = MbonReadout.load(args.mbon)
|
||||
print(f"считывание MBON: {args.mbon}")
|
||||
over = {}
|
||||
if args.mbon_blend is not None:
|
||||
|
|
@ -209,23 +221,28 @@ def main() -> None:
|
|||
params = Params(**over)
|
||||
laterals = tuple(float(x) for x in args.laterals.split(","))
|
||||
|
||||
all_rec = []
|
||||
tasks = []
|
||||
for p in find_bags(args.root):
|
||||
if p.name == HOLDOUT:
|
||||
continue
|
||||
t0 = time.time()
|
||||
rd = readout
|
||||
if folds:
|
||||
rd = folds.get(p.name)
|
||||
if rd is None:
|
||||
rd = args.mbon
|
||||
if fold_paths:
|
||||
rd = fold_paths.get(p.name, "")
|
||||
if not rd:
|
||||
print(f" внимание: для {p.name} нет своей складки — пропуск")
|
||||
continue
|
||||
rec = run_bag(p, params, memory, args.limit, args.d_start, laterals,
|
||||
args.seed, readout=rd)
|
||||
all_rec.extend(rec)
|
||||
tasks.append((p, params, args.limit, args.d_start, laterals, args.seed,
|
||||
args.memory or "", rd))
|
||||
|
||||
# Печатается по готовности, собирается по номеру задачи: порядок сценариев
|
||||
# в файле не должен зависеть от того, какой бэг досчитался первым.
|
||||
slots: list = [None] * len(tasks)
|
||||
for i, task, rec, secs in P.run(_work, tasks, args.jobs):
|
||||
slots[i] = rec
|
||||
n = sum(len(r["d"]) for r in rec)
|
||||
print(f" {p.name:42s} сценариев {len(rec):3d}, наблюдений {n:6d}, "
|
||||
f"{time.time()-t0:6.1f} с", flush=True)
|
||||
print(f" {task[0].name:42s} сценариев {len(rec):3d}, наблюдений {n:6d}, "
|
||||
f"{secs:6.1f} с", flush=True)
|
||||
all_rec = [r for rec in slots if rec for r in rec]
|
||||
|
||||
with open(args.out, "w", encoding="utf-8") as f:
|
||||
json.dump(all_rec, f, ensure_ascii=False)
|
||||
|
|
|
|||
|
|
@ -27,11 +27,11 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import time
|
||||
|
||||
import numpy as np
|
||||
|
||||
import _bootstrap as B # noqa: F401
|
||||
import _parallel as P
|
||||
from flyguard.bag import Bag, find_bags
|
||||
from flyguard.mushroom_body import FEATURES, describe
|
||||
from flyguard.pipeline import FlyGuard, Params
|
||||
|
|
@ -41,6 +41,12 @@ HOLDOUT = "doubleT_obstacle" # там реальный объект — т
|
|||
MIN_OVERLAP = 0.5 # доля лучей ядра, пришедших от предмета
|
||||
|
||||
|
||||
def _work(task):
|
||||
"""Одна задача — один бэг. Верхнего уровня: иначе не передать в процесс."""
|
||||
path, params, limit, d_starts, laterals, seed = task
|
||||
return collect_bag(path, params, limit, d_starts, laterals, seed)
|
||||
|
||||
|
||||
def ego_track(bag: Bag, params: Params, limit: int):
|
||||
"""Первый проход: пройденный путь на каждом кадре и общая решётка."""
|
||||
fg = FlyGuard(params, memory=None)
|
||||
|
|
@ -134,6 +140,7 @@ def main() -> None:
|
|||
# красивая и бессмысленная, а в тоннеле у оси полно штатных конструкций.
|
||||
ap.add_argument("--laterals", default="0,-0.6,0.6,-1.2,1.2")
|
||||
ap.add_argument("--seed", type=int, default=20260921)
|
||||
P.add_argument(ap)
|
||||
args = ap.parse_args()
|
||||
|
||||
d_starts = tuple(float(x) for x in args.d_starts.split(","))
|
||||
|
|
@ -141,20 +148,22 @@ def main() -> None:
|
|||
params = Params()
|
||||
B.CACHE.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
parts = {}
|
||||
for p in find_bags(args.root):
|
||||
if p.name == HOLDOUT:
|
||||
continue
|
||||
t0 = time.time()
|
||||
got = collect_bag(p, params, args.limit, d_starts, laterals, args.seed)
|
||||
bags = [p for p in find_bags(args.root) if p.name != HOLDOUT]
|
||||
tasks = [(p, params, args.limit, d_starts, laterals, args.seed) for p in bags]
|
||||
# Печатается по готовности, складывается по номеру задачи: порядок бэгов
|
||||
# в файле не должен зависеть от того, какой из них досчитался первым.
|
||||
slots: list = [None] * len(tasks)
|
||||
for i, task, got, secs in P.run(_work, tasks, args.jobs):
|
||||
name = task[0].name
|
||||
if got is None:
|
||||
print(f" {p.name:42s} пропущен")
|
||||
print(f" {name:42s} пропущен")
|
||||
continue
|
||||
parts[p.name] = got
|
||||
slots[i] = (name, got)
|
||||
Xb, yb = got[0], got[1]
|
||||
print(f" {p.name:42s} {Xb.shape[0]:7d} кандидатов, "
|
||||
print(f" {name:42s} {Xb.shape[0]:7d} кандидатов, "
|
||||
f"предметов {int(yb.sum()):6d} ({yb.mean():5.1%}), "
|
||||
f"{time.time() - t0:6.0f} с", flush=True)
|
||||
f"{secs:6.0f} с", flush=True)
|
||||
parts = {name: got for name, got in (s for s in slots if s is not None)}
|
||||
|
||||
if not parts:
|
||||
raise SystemExit("ничего не собрано")
|
||||
|
|
|
|||
Loading…
Reference in a new issue