Хранилище: писатель пачками + выживание при SIGTERM — конец потерь записей

GPU-эволюция (380+/с) обгоняла писателя (~150/с, коммит=fsync на каждую
запись): хвост очереди копился в RAM минутами и ТЕРЯЛСЯ при рестарте
сервиса (~80к записей за один рестарт) — нарушение «сохраняем всё».
- run_writer_process: вставка пачками до 500 в одной транзакции; при ошибке
  пачки — откат и поштучный поиск виновника; SIGTERM игнорирует (дописывает
  всё до None)
- run_evolution: SIGTERM -> исключение -> штатный finally досылает None и
  ждёт писателя; в лог явные строки про дописывание базы и полировку
  (чтобы «=== готово ===» не выглядело остановкой)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
jze9
2026-07-08 04:18:11 +05:00
parent ac459d26a4
commit e6b877bab7
2 changed files with 52 additions and 11 deletions

View File

@@ -11,6 +11,7 @@ import contextlib
import copy import copy
import multiprocessing import multiprocessing
import random import random
import signal
from concurrent.futures import ProcessPoolExecutor from concurrent.futures import ProcessPoolExecutor
from pathlib import Path from pathlib import Path
@@ -123,6 +124,14 @@ def run_evolution(
writer = ctx.Process(target=run_writer_process, args=(queue, db_path)) writer = ctx.Process(target=run_writer_process, args=(queue, db_path))
writer.start() writer.start()
# SIGTERM (остановка/рестарт systemd) -> исключение -> блок finally ниже
# штатно досылает None и ЖДЁТ писателя: всё уже оценённое дописывается в
# базу, а не теряется вместе с очередью в RAM («сохраняем всё»).
def _terminate(_sig, _frm):
raise KeyboardInterrupt("SIGTERM: дописываем очередь и выходим")
signal.signal(signal.SIGTERM, _terminate)
population = [sample_genome(db, bounds, rng) for _ in range(population_size)] population = [sample_genome(db, bounds, rng) for _ in range(population_size)]
best_pair: tuple[Genome, EvaluationResult] | None = None best_pair: tuple[Genome, EvaluationResult] | None = None
n_evaluated = 0 n_evaluated = 0
@@ -166,14 +175,18 @@ def run_evolution(
population = next_population population = next_population
finally: finally:
logger.finish() logger.finish()
logger._write("дописываю хвост очереди в базу (писатель)…")
queue.put(None) queue.put(None)
writer.join() writer.join()
logger._write("база дописана полностью")
assert best_pair is not None assert best_pair is not None
best_genome, best_result = best_pair best_genome, best_result = best_pair
if polish and best_result.feasible: if polish and best_result.feasible:
logger._write("полировка лучшего генома (Nelder-Mead, CPU) — это НЕ остановка, новый цикл стартует после неё")
polished_genome, polished_result = polish_best(best_genome, db, bounds) polished_genome, polished_result = polish_best(best_genome, db, bounds)
logger._write(f"полировка готова: КПД {best_result.fitness*100:.2f}% -> {polished_result.fitness*100:.2f}%")
if polished_result.fitness > best_result.fitness: if polished_result.fitness > best_result.fitness:
conn_queue = ctx.Queue() conn_queue = ctx.Queue()
polish_writer = ctx.Process(target=run_writer_process, args=(conn_queue, db_path)) polish_writer = ctx.Process(target=run_writer_process, args=(conn_queue, db_path))

View File

@@ -107,20 +107,48 @@ def run_writer_process(queue, db_path: Path) -> None:
единственный writer на файл базы, пока воркеры sweep/evolve только единственный writer на файл базы, пока воркеры sweep/evolve только
кладут результаты в очередь. кладут результаты в очередь.
Один "плохой" прогон (например, коллизия run_id) не должен уронить ПАЧКАМИ, не по одной: коммит SQLite = fsync (~5-10мс), и по-одному
писатель и остановить осушение очереди на много часов sweep — при писатель выдавал ~150/с — GPU-эволюция (380+/с) копила хвост очереди в
ошибке вставки конкретная запись пропускается с сообщением в stderr, RAM на минуты, который терялся при рестарте сервиса («сохраняем всё»
а не молча и не ценой падения всего процесса. нарушалось). Теперь копим до BATCH записей (или пока очередь не
опустела) и вставляем одной транзакцией — тысячи вставок в секунду,
хвост исчезает.
SIGTERM игнорируется: при остановке сервиса systemd шлёт TERM всей
группе; писатель обязан дожить до None от главного процесса и дописать
ВСЁ, что уже оценено, а не умереть с полной очередью.
Одна "плохая" пачка (например, коллизия run_id) не роняет писатель:
откатываемся и вставляем пачку по одной, пропуская только виновника.
""" """
import queue as queue_mod
import signal
signal.signal(signal.SIGTERM, signal.SIG_IGN)
BATCH = 500
conn = open_connection(db_path) conn = open_connection(db_path)
try: try:
while True: finished = False
record = queue.get() while not finished:
if record is None: batch = [queue.get()] # блокируемся на первой
break while len(batch) < BATCH:
try:
batch.append(queue.get_nowait())
except queue_mod.Empty:
break
if batch[-1] is None:
finished = True
batch.pop()
if not batch:
continue
try: try:
insert_run(conn, record) insert_runs(conn, batch)
except sqlite3.Error as exc: except sqlite3.Error:
print(f"[storage] не удалось записать run_id={record.run_id}: {exc}", file=sys.stderr) conn.rollback()
for record in batch: # ищем виновника, остальное сохраняем
try:
insert_run(conn, record)
except sqlite3.Error as exc:
print(f"[storage] не удалось записать run_id={record.run_id}: {exc}", file=sys.stderr)
finally: finally:
conn.close() conn.close()