From 7732577f31a6e6a09647308c23d6ffe2e9e997bc Mon Sep 17 00:00:00 2001 From: jze9 Date: Sun, 12 Jul 2026 22:29:17 +0500 Subject: [PATCH] =?UTF-8?q?=D0=AD=D0=B2=D0=BE=D0=BB=D1=8E=D1=86=D0=B8?= =?UTF-8?q?=D1=8F:=20=D0=B2=20=D0=B1=D0=B0=D0=B7=D1=83=20=D1=82=D0=BE?= =?UTF-8?q?=D0=BB=D1=8C=D0=BA=D0=BE=20=D1=83=D0=BB=D1=83=D1=87=D1=88=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D1=8F=20=D0=BB=D1=83=D1=87=D1=88=D0=B8=D1=85=20+?= =?UTF-8?q?=20=D1=81=D1=87=D1=91=D1=82=D1=87=D0=B8=D0=BA=D0=B8=20=D0=BF?= =?UTF-8?q?=D0=BE=20=D0=B2=D1=81=D0=B5=D0=BC;=20=D0=B2=D1=8B=D0=B1=D0=BE?= =?UTF-8?q?=D1=80=D0=BA=D0=B0=2010000?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit По требованию пользователя «не сохранять всё подряд, а только то, что нужно для эволюции»: рядовые геномы поколения больше НЕ пишутся строками в runs — писатель получает лёгкие StatsOnly-инкременты (счётчики, гистограмма, причины отказов на дашборде остаются честными по ВСЕМ оценкам), а полная строка пишется только когда лучший геном острова улучшился (+ полировка). Вставка миллионов строк в HDD-базу с 6 индексами была узким местом конвейера: GPU считал цикл ~23 мин, дозапись хвоста очереди шла часами (~150 строк/с), реальная скорость была ~165 оценок/с вместо ~1300/с. gpu_record_worker стал не нужен — удалён. Выборка эволюции уменьшена до 10000 на поколение (8 островов x 1250, было 8x4000=32000) — по просьбе пользователя, чтобы поколения сменялись быстрее; поколений за цикл теперь 120 (было 60) — при освобождённом писателе цикл углубляется вдвое дальше за то же время. --- deploy/evolve_forever.sh | 6 +++- src/gausse/gpu/batch_sweep.py | 37 -------------------- src/gausse/optim/evolutionary.py | 59 ++++++++++++++++++++++---------- src/gausse/storage/database.py | 53 ++++++++++++++++++++-------- src/gausse/storage/schema.py | 19 +++++++--- tests/test_evolutionary.py | 30 +++++++++++++--- tests/test_storage_database.py | 30 ++++++++++++++++ 7 files changed, 153 insertions(+), 81 deletions(-) diff --git a/deploy/evolve_forever.sh b/deploy/evolve_forever.sh index 02115f4..4ac08d9 100755 --- a/deploy/evolve_forever.sh +++ b/deploy/evolve_forever.sh @@ -21,10 +21,14 @@ except Exception: PY ) echo "$(date -Is) ═════ ЦИКЛ #$cycle стартует (конвейер бесконечный; лучший КПД за всё время: $best) ═════" >> "$LOG" + # выборка 10000 на поколение (8 островов x 1250) — по просьбе пользователя: + # меньше популяция + больше поколений за цикл = быстрее углубляется. + # Рядовые оценки в базу не пишутся (только счётчики) — писатель больше + # не душит конвейер, GPU занят почти всё время цикла. python3 -c "from pathlib import Path from gausse.optim.evolutionary import run_evolution try: - run_evolution(Path('results/gausse.sqlite3'), n_generations=60, population_size=4000, seed=None, polish=True, use_gpu=True, islands=8) + run_evolution(Path('results/gausse.sqlite3'), n_generations=120, population_size=1250, seed=None, polish=True, use_gpu=True, islands=8) except KeyboardInterrupt: pass # штатная остановка сервиса: база уже дописана в finally" echo "$(date -Is) ═════ ЦИКЛ #$cycle завершён — это НЕ остановка: через 3с стартует цикл #$((cycle+1)) со свежей случайной популяцией ═════" >> "$LOG" diff --git a/src/gausse/gpu/batch_sweep.py b/src/gausse/gpu/batch_sweep.py index 03212d3..8c24c05 100644 --- a/src/gausse/gpu/batch_sweep.py +++ b/src/gausse/gpu/batch_sweep.py @@ -299,43 +299,6 @@ def evaluate_genomes_gpu( return results -def gpu_record_worker(payload) -> RunRecord: - """Собирает RunRecord (decode + стоимость + JSON) в процессе пула — - параллельно с GPU-оценкой следующих батчей. - - detail НЕ строится (build_detail дорогая — полная физика по каждой - ступени — и не нужна для хранения): genome_json уже достаточен, чтобы - воспроизвести всё через decode()+build_detail() на лету при просмотре - конкретного прогона (см. web/stats.run_detail). Пишем дешёвую - decoded_summary() — иначе на миллионах прогонов это дорогая работа, - съедающая и CPU (душит GPU), и диск (десятки ГБ на пустом месте). - - Работает в воркере ProcessPoolExecutor с worker_context.init_worker - (база компонентов уже загружена в процессе). - """ - from gausse.optim import worker_context - from gausse.optim.objective import MODEL_VERSION, compute_cost_rub, decoded_summary - - genome, outcome, search_mode = payload - db, bounds = worker_context.db, worker_context.bounds - config, _ix, _iv = decode(genome, db, bounds) - return RunRecord( - run_id=str(uuid.uuid4()), - timestamp=datetime.now(timezone.utc).isoformat(), - search_mode=search_mode, - genome_json=json.dumps(genome_to_dict(genome)), - decoded_summary_json=json.dumps(decoded_summary(genome, db), ensure_ascii=False), - feasible=outcome["feasible"], - model_version=MODEL_VERSION, - infeasible_reason=outcome["reason"], - failed_stage_index=outcome["failed_stage_index"], - efficiency=outcome["efficiency"], - exit_velocity_mps=outcome["exit_velocity_mps"], - cost_rub=compute_cost_rub(config, db), - energy_breakdown_json=json.dumps(outcome["energy_breakdown"]) if outcome["energy_breakdown"] else None, - ) - - def run_gpu_sweep( db_path: Path, n_runs: int, diff --git a/src/gausse/optim/evolutionary.py b/src/gausse/optim/evolutionary.py index e86a220..2a6b08a 100644 --- a/src/gausse/optim/evolutionary.py +++ b/src/gausse/optim/evolutionary.py @@ -19,10 +19,10 @@ from scipy.optimize import minimize from gausse.components.database import ComponentDatabase from gausse.optim import worker_context -from gausse.optim.objective import EvaluationResult, build_run_record, evaluate +from gausse.optim.objective import EvaluationResult, build_run_record, compute_cost_rub, evaluate from gausse.optim.progress_log import ProgressLogger, default_log_path -from gausse.optim.search_space import Genome, SearchBounds, crossover, mutate, repair, sample_genome -from gausse.storage.database import insert_run, open_connection, run_writer_process +from gausse.optim.search_space import Genome, SearchBounds, crossover, decode, mutate, repair, sample_genome +from gausse.storage.database import StatsOnly, insert_run, open_connection, run_writer_process from gausse.storage.schema import RunRecord @@ -31,6 +31,15 @@ def _evaluate_genome(genome: Genome) -> tuple[Genome, EvaluationResult]: return genome, result +def _best_record(genome: Genome, result: EvaluationResult, db, bounds, mode: str) -> RunRecord: + """Полная строка БД для сохраняемого лучшего генома. GPU-путь не считает + стоимость для рядовых оценок — дочитываем её только здесь, для избранных.""" + if result.cost_rub is None: + config, _, _ = decode(genome, db, bounds) + result.cost_rub = compute_cost_rub(config, db) + return build_run_record(genome, result, db, search_mode=mode) + + def _tournament_select( pairs: list[tuple[Genome, EvaluationResult]], rng: random.Random, k: int = 3 ) -> Genome: @@ -162,33 +171,24 @@ def run_evolution( initializer=worker_context.init_worker, initargs=(data_dir, bounds), ) as executor: + # фитнес лучшей УЖЕ ЗАПИСАННОЙ строки каждого острова: полные + # строки в базу идут только когда остров улучшился — иначе элитизм + # дублировал бы один и тот же геном в топ дашборда каждое поколение + written_best_fitness = [float("-inf")] * islands + for generation in range(n_generations): # все острова — ОДНИМ плоским батчем (жирный батч кормит GPU) flat = [g for island in populations for g in island] if gpu_xp is not None: - from gausse.gpu.batch_sweep import evaluate_genomes_gpu, gpu_record_worker + from gausse.gpu.batch_sweep import evaluate_genomes_gpu results = evaluate_genomes_gpu(gpu_xp, flat, db, bounds, executor=executor) pairs_flat = list(zip(flat, results)) - payload = [ - (g, { - "feasible": r.feasible, - "efficiency": r.efficiency, - "exit_velocity_mps": r.exit_velocity_mps, - "reason": r.reason, - "failed_stage_index": r.failed_stage_index, - "energy_breakdown": r.energy_breakdown, - }, mode) - for g, r in pairs_flat - ] - records = executor.map(gpu_record_worker, payload, chunksize=64) else: pairs_flat = list(executor.map(_evaluate_genome, flat)) - records = (build_run_record(g, r, db, search_mode=mode) for g, r in pairs_flat) n_evaluated += len(pairs_flat) - for record, (genome, result) in zip(records, pairs_flat): - queue.put(record) + for genome, result in pairs_flat: logger.update(result.feasible, result.efficiency) if best_pair is None or result.fitness > best_pair[1].fitness: best_pair = (genome, result) @@ -196,10 +196,19 @@ def run_evolution( # селекция и скрещивание — НА КАЖДОМ ОСТРОВЕ независимо island_bests: list[Genome] = [] next_populations = [] + recorded = set() # id() результатов, ушедших полной строкой for i in range(islands): pairs = pairs_flat[i * population_size:(i + 1) * population_size] pairs.sort(key=lambda pair: pair[1].fitness, reverse=True) island_bests.append(pairs[0][0]) + # улучшение лучшего на острове — полная строка в базу + # (топы дашборда, best за всё время, отчёты). Новый + # глобальный best — всегда чей-то островной best, так что + # он не теряется. + if pairs[0][1].fitness > written_best_fitness[i]: + written_best_fitness[i] = pairs[0][1].fitness + queue.put(_best_record(pairs[0][0], pairs[0][1], db, bounds, mode)) + recorded.add(id(pairs[0][1])) next_population = [copy.deepcopy(pair[0]) for pair in pairs[:elitism]] while len(next_population) < population_size: parent_a = _tournament_select(pairs, rng) @@ -208,6 +217,18 @@ def run_evolution( child = mutate(child, db, bounds, rng, rate=mutation_rate) next_population.append(child) next_populations.append(next_population) + + # рядовые оценки — в базу ТОЛЬКО агрегатами (счётчики/ + # гистограмма честно считают все оценки, вставка полной строки + # выше сама инкрементирует их для избранных — без двойного + # счёта). Хранить миллионы проходных геномов незачем + # («сохраняй то, что нужно для эволюции»), а вставка каждого + # в HDD-базу с 6 индексами была узким местом конвейера: + # GPU считал цикл ~23 мин, дозапись хвоста шла часами. + for _genome, result in pairs_flat: + if id(result) not in recorded: + queue.put(StatsOnly(mode, result.feasible, result.efficiency, result.reason)) + # миграция по кольцу: лучший острова i замещает слот у острова i+1 if islands > 1 and (generation + 1) % migrate_every == 0: for i in range(islands): diff --git a/src/gausse/storage/database.py b/src/gausse/storage/database.py index 57f695f..5a6c1e4 100644 --- a/src/gausse/storage/database.py +++ b/src/gausse/storage/database.py @@ -9,7 +9,7 @@ import sqlite3 import sys -from dataclasses import asdict, fields +from dataclasses import asdict, dataclass, fields from pathlib import Path from gausse.storage.schema import CREATE_INDEXES_SQL, CREATE_RUNS_TABLE_SQL, CREATE_STATS_SQL, RunRecord @@ -17,6 +17,25 @@ from gausse.storage.schema import CREATE_INDEXES_SQL, CREATE_RUNS_TABLE_SQL, CRE _COLUMNS = [f.name for f in fields(RunRecord)] +@dataclass(frozen=True) +class StatsOnly: + """Лёгкое сообщение писателю: учесть оценку в агрегатах дашборда + (счётчики/гистограмма/причины), НЕ сохраняя строку в runs. + + Эволюция шлёт такое для РЯДОВЫХ геномов поколения, а полной строкой + пишет только улучшения лучших (см. optim/evolutionary.py): хранить + миллионы проходных геномов незачем, а вставка каждого в HDD-базу с + 6 индексами была узким местом всего конвейера (GPU считал цикл за + ~23 мин, дозапись хвоста очереди шла часами при ~150 строк/с). + Атрибуты названы как у RunRecord — _apply_stats работает с обоими. + """ + + search_mode: str + feasible: bool + efficiency: float | None + infeasible_reason: str | None + + def open_connection(db_path: Path) -> sqlite3.Connection: db_path = Path(db_path) db_path.parent.mkdir(parents=True, exist_ok=True) @@ -48,8 +67,10 @@ def _efficiency_bucket(efficiency: float | None) -> int | None: return b if 0 <= b <= 9 else None -def _apply_stats(conn: sqlite3.Connection, records: list[RunRecord]) -> None: - """Инкремент агрегатов дашборда — вызывается внутри транзакции вставки.""" +def _apply_stats(conn: sqlite3.Connection, records: list) -> None: + """Инкремент агрегатов дашборда — вызывается внутри транзакции вставки. + + Принимает и RunRecord, и StatsOnly (одинаковые имена атрибутов).""" n_feasible = sum(1 for r in records if r.feasible) conn.execute( "INSERT INTO stats_counters(key, value) VALUES('total', ?) " @@ -140,17 +161,20 @@ def insert_run(conn: sqlite3.Connection, record: RunRecord) -> None: insert_runs(conn, [record]) -def insert_runs(conn: sqlite3.Connection, records: list[RunRecord]) -> None: +def insert_runs(conn: sqlite3.Connection, records: list) -> None: + """Пишет пачку RunRecord/StatsOnly одной транзакцией: строки runs — только + для RunRecord, агрегаты дашборда — по ВСЕМ элементам пачки (либо всё, + либо ничего — иначе цифры на морде разъедутся с таблицей).""" if not records: return - placeholders = ", ".join(f":{name}" for name in _COLUMNS) - columns = ", ".join(_COLUMNS) - conn.executemany( - f"INSERT INTO runs ({columns}) VALUES ({placeholders})", - [_record_to_params(r) for r in records], - ) - # агрегаты дашборда — в той же транзакции: либо и строки, и счётчики, - # либо ничего (иначе цифры на морде разъедутся с таблицей) + rows = [r for r in records if isinstance(r, RunRecord)] + if rows: + placeholders = ", ".join(f":{name}" for name in _COLUMNS) + columns = ", ".join(_COLUMNS) + conn.executemany( + f"INSERT INTO runs ({columns}) VALUES ({placeholders})", + [_record_to_params(r) for r in rows], + ) _apply_stats(conn, records) conn.commit() @@ -208,7 +232,7 @@ def fetch_runs( def run_writer_process(queue, db_path: Path) -> None: - """Читает `RunRecord` из очереди и пишет в SQLite, пока не придёт None. + """Читает `RunRecord`/`StatsOnly` из очереди и пишет в SQLite до None. Предназначен для запуска как отдельный `multiprocessing.Process` — единственный writer на файл базы, пока воркеры sweep/evolve только @@ -262,7 +286,8 @@ def run_writer_process(queue, db_path: Path) -> None: try: insert_run(conn, record) except sqlite3.Error as exc: - print(f"[storage] не удалось записать run_id={record.run_id}: {exc}", file=sys.stderr) + rid = getattr(record, "run_id", "") + print(f"[storage] не удалось записать run_id={rid}: {exc}", file=sys.stderr) # авто-чекпоинт голодает при постоянном потоке — двигаем WAL сами, # PASSIVE не блокирует читателей if time_mod.monotonic() - last_checkpoint > CHECKPOINT_EVERY_S: diff --git a/src/gausse/storage/schema.py b/src/gausse/storage/schema.py index 8959465..09ca8f0 100644 --- a/src/gausse/storage/schema.py +++ b/src/gausse/storage/schema.py @@ -1,9 +1,18 @@ -"""Схема таблицы `runs` — ВСЕ прогоны (успешные и нет), без прикрас. +"""Схема таблицы `runs` + агрегаты дашборда. -Нереализуемые конфигурации хранятся наравне с успешными: `feasible=0` + -`infeasible_reason` с честной причиной. Ничего не отбрасывается перед -записью — отбор "хороших" результатов делается запросом к этой таблице -постфактум, а не фильтрацией на входе. +Честность цифр — через агрегаты: счётчики/гистограмма/причины (stats_*) +инкрементируются для КАЖДОЙ оценки, включая нереализуемые, в той же +транзакции писателя. А вот полные СТРОКИ в `runs` пишутся по-разному: + +- sweep/manual: каждая оценка — строка (исследовательский режим); +- evolve (с 2026-07-12): строка — только когда лучший геном острова + УЛУЧШИЛСЯ (+ полировка); рядовые миллионы геномов поколения идут в базу + как StatsOnly-инкременты счётчиков. Хранить их незачем («сохраняй то, + что нужно для эволюции»), а вставка каждого в HDD-базу с 6 индексами + была узким местом: GPU считал цикл ~23 мин, дозапись шла часами. + +Нереализуемые не приукрашиваются: у строк feasible=0 + честная причина, +у агрегатов — счётчик причин отказов. """ from dataclasses import dataclass diff --git a/tests/test_evolutionary.py b/tests/test_evolutionary.py index 9d18c26..f73cc8d 100644 --- a/tests/test_evolutionary.py +++ b/tests/test_evolutionary.py @@ -8,7 +8,14 @@ from gausse.storage.database import count_runs, fetch_runs, open_connection DB = ComponentDatabase.load() -def test_run_evolution_writes_all_evaluations_and_returns_summary(tmp_path): +def _stats_total(conn) -> int: + row = conn.execute("SELECT value FROM stats_counters WHERE key='total'").fetchone() + return row[0] if row else 0 + + +def test_run_evolution_counts_all_evaluations_but_stores_only_best_rows(tmp_path): + """Контракт хранения evolve: агрегаты честно считают ВСЕ оценки, + а полные строки — только улучшения лучшего генома острова.""" db_path = tmp_path / "runs.sqlite3" bounds = SearchBounds(max_stages=2) @@ -26,9 +33,16 @@ def test_run_evolution_writes_all_evaluations_and_returns_summary(tmp_path): assert "best_fitness" in summary conn = open_connection(db_path) - assert count_runs(conn) == 12 + assert _stats_total(conn) == 12 # счётчики дашборда видят все оценки + n_rows = count_runs(conn) + # строк меньше, чем оценок: максимум 1 улучшение на остров на поколение + assert 1 <= n_rows <= 2 rows = fetch_runs(conn) assert all(r.search_mode == "evolve" for r in rows) + # лучший найденный обязан лежать в базе полной строкой + best_row_fitness = max((r.efficiency if r.feasible else -1.0) for r in rows) + if summary["best_feasible"]: + assert best_row_fitness == summary["best_efficiency"] conn.close() @@ -67,7 +81,7 @@ def test_polish_does_not_make_the_best_genome_worse(): def test_run_evolution_gpu_batch_mode(tmp_path): """GPU-режим эволюции (батч-оценка поколения; здесь numpy-бэкенд): - пишет прогоны в ту же БД и находит реализуемые конфигурации.""" + агрегаты считают все оценки, строки — только улучшения лучшего.""" from gausse.storage.database import count_runs, fetch_runs, open_connection db_path = tmp_path / "evo_gpu.sqlite3" @@ -76,9 +90,13 @@ def test_run_evolution_gpu_batch_mode(tmp_path): ) assert summary["n_evaluated"] == 120 conn = open_connection(db_path) - assert count_runs(conn) >= 120 + assert _stats_total(conn) == 120 + assert 1 <= count_runs(conn) <= 3 # 1 остров x 3 поколения, только улучшения rows = fetch_runs(conn, limit=5) assert all(r.search_mode.startswith("evolve-gpu") for r in rows) + # у сохранённых лучших должна быть посчитана стоимость (GPU-путь + # не считает её для рядовых, но дочитывает для избранных) + assert all(r.cost_rub is not None for r in rows) def test_run_evolution_islands(tmp_path): @@ -92,4 +110,6 @@ def test_run_evolution_islands(tmp_path): ) assert summary["n_evaluated"] == 4 * 15 * 3 conn = open_connection(db_path) - assert count_runs(conn) == 4 * 15 * 3 + assert _stats_total(conn) == 4 * 15 * 3 + # максимум одно улучшение на остров на поколение + assert 3 <= count_runs(conn) <= 4 * 3 diff --git a/tests/test_storage_database.py b/tests/test_storage_database.py index e4d6fa6..ce19e54 100644 --- a/tests/test_storage_database.py +++ b/tests/test_storage_database.py @@ -99,6 +99,36 @@ def test_writer_survives_duplicate_run_id_and_keeps_draining_queue(tmp_path): conn.close() +def test_writer_stats_only_bumps_counters_without_rows(tmp_path): + """StatsOnly-сообщения (эволюция шлёт их для рядовых геномов): агрегаты + дашборда растут, а строк в runs не появляется.""" + from gausse.storage.database import StatsOnly + + db_path = tmp_path / "runs.sqlite3" + ctx = multiprocessing.get_context("spawn") + queue = ctx.Queue() + + writer = ctx.Process(target=run_writer_process, args=(queue, db_path)) + writer.start() + + queue.put(_record("best-one")) # одна полная строка + queue.put(StatsOnly("evolve", True, 0.42, None)) + queue.put(StatsOnly("evolve", False, None, "датчик: сигнал ниже порога")) + queue.put(StatsOnly("evolve", False, None, None)) + queue.put(None) + writer.join(timeout=10) + + assert not writer.is_alive() + conn = open_connection(db_path) + assert count_runs(conn) == 1 # строка только у полной записи + counters = dict(conn.execute("SELECT key, value FROM stats_counters")) + assert counters["total"] == 4 # но посчитаны все четыре + assert counters["feasible"] == 2 + reasons = dict(conn.execute("SELECT reason, count FROM stats_reasons")) + assert reasons.get("датчик") == 1 + conn.close() + + def _writer_worker(queue, worker_id, n): for i in range(n): queue.put(_record(f"proc-{worker_id}-run-{i}"))