Эволюция: в базу только улучшения лучших + счётчики по всем; выборка 10000
По требованию пользователя «не сохранять всё подряд, а только то, что нужно для эволюции»: рядовые геномы поколения больше НЕ пишутся строками в runs — писатель получает лёгкие StatsOnly-инкременты (счётчики, гистограмма, причины отказов на дашборде остаются честными по ВСЕМ оценкам), а полная строка пишется только когда лучший геном острова улучшился (+ полировка). Вставка миллионов строк в HDD-базу с 6 индексами была узким местом конвейера: GPU считал цикл ~23 мин, дозапись хвоста очереди шла часами (~150 строк/с), реальная скорость была ~165 оценок/с вместо ~1300/с. gpu_record_worker стал не нужен — удалён. Выборка эволюции уменьшена до 10000 на поколение (8 островов x 1250, было 8x4000=32000) — по просьбе пользователя, чтобы поколения сменялись быстрее; поколений за цикл теперь 120 (было 60) — при освобождённом писателе цикл углубляется вдвое дальше за то же время.
This commit is contained in:
@@ -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"
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
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 records],
|
||||
[_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", "<stats-only>")
|
||||
print(f"[storage] не удалось записать run_id={rid}: {exc}", file=sys.stderr)
|
||||
# авто-чекпоинт голодает при постоянном потоке — двигаем WAL сами,
|
||||
# PASSIVE не блокирует читателей
|
||||
if time_mod.monotonic() - last_checkpoint > CHECKPOINT_EVERY_S:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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}"))
|
||||
|
||||
Reference in New Issue
Block a user