# Boucle d'entrainement avec checkpoint, thermal watchdog et AMP
# Utilise le module shared/gpu_training.py pour eviter de dupliquer le code
import sys
import importlib
# Resolve shared module: works from notebook dir, Papermill, and CLI
_shared_candidates = [
os.path.abspath(os.path.join(os.getcwd(), '..')), # Running from Python/ dir
os.path.abspath(os.path.join(os.getcwd(), 'MyIA.AI.Notebooks', 'QuantConnect')), # Running from repo root
os.path.abspath(os.path.join(os.getcwd(), 'QuantConnect')), # Running from notebooks root
]
for _p in _shared_candidates:
if os.path.isfile(os.path.join(_p, 'shared', 'gpu_training.py')):
if _p not in sys.path:
sys.path.insert(0, _p)
break
from shared.gpu_training import TrainingCheckpoint, setup_amp, get_gpu_temp
# --- Notebook directory (robust across Jupyter / Papermill / CLI) ---
_nb_dir = os.path.abspath(os.path.join(os.getcwd(), 'MyIA.AI.Notebooks', 'QuantConnect', 'Python'))
if not os.path.isfile(os.path.join(_nb_dir, 'QC-Py-31-Transformer-Training.ipynb')):
_nb_dir = os.getcwd() # Already in notebook dir (Jupyter)
# --- Configuration checkpoint ---
checkpoint_path = os.path.join(_nb_dir, 'transformer_checkpoint.pt')
model_save_path = os.path.join(_nb_dir, 'transformer_multiasset_model.pt')
# --- AMP (Mixed Precision) ---
use_amp, grad_scaler = setup_amp()
print(f"AMP: {'Active (GPU)' if use_amp else 'Desactive (CPU)'}")
# --- Optimiseur et scheduler ---
optimizer = optim.AdamW(
model.parameters(),
lr=LR,
weight_decay=0.01
)
# Warmup lineaire avant cosine. Controle du 26/09 : sans warmup, le LR
# plein des l'epoch 1 (AdamW) aplatit l'attention a l'uniforme -- poids
# 1/60 = 0.017 sur toutes les tetes mesure au run du 26/09 -- l'encodeur
# ne distingue plus les pas de temps et les deux tetes sortent des
# predictions constantes (Correlation nan, DirAcc = classe majoritaire).
warmup_sched = optim.lr_scheduler.LinearLR(
optimizer,
start_factor=0.01,
end_factor=1.0,
total_iters=WARMUP_EPOCHS
)
cosine_sched = optim.lr_scheduler.CosineAnnealingLR(
optimizer,
T_max=EPOCHS - WARMUP_EPOCHS,
eta_min=LR * 0.01
)
scheduler = optim.lr_scheduler.SequentialLR(
optimizer,
[warmup_sched, cosine_sched],
milestones=[WARMUP_EPOCHS]
)
# --- Historique ---
history = {
'train_loss': [], 'val_loss': [],
'train_reg_loss': [], 'val_reg_loss': [],
'train_cls_loss': [], 'val_cls_loss': [],
'train_dir_acc': [], 'val_dir_acc': []
}
# --- Resume checkpoint (3 cas: modele final / checkpoint / from scratch) ---
# FORCE_RETRAIN : execution de reference depuis zero (#17516) -- un .pt residuel
# sur la machine ne doit jamais produire une courbe d'apprentissage vide.
ckpt_manager = TrainingCheckpoint(
checkpoint_path=checkpoint_path,
model_save_path=model_save_path,
max_temp=80,
thermal_check_every=10,
cool_sleep=15,
patience=7
)
if FORCE_RETRAIN:
print("FORCE_RETRAIN=True : checkpoint/.pt ignores, entrainement depuis zero.")
start_epoch = 0
else:
start_epoch, history = ckpt_manager.resume(
model, optimizer, scheduler, grad_scaler,
device=device, default_history=history
)
# --- Boucle d'entrainement ---
if start_epoch >= 0 and start_epoch < EPOCHS:
print("=" * 70)
print("ENTRAINEMENT TRANSFORMER ENCODER (AMP + Thermal Watchdog)")
print("=" * 70)
print(f"Device: {device}")
print(f"Epochs: {start_epoch+1}-{EPOCHS}, Batch: {BATCH_SIZE}, LR: {LR}")
print(f"Thermal limit: {ckpt_manager.max_temp}C")
print()
for epoch in range(start_epoch, EPOCHS):
epoch_start = time.time()
# Thermal check via le helper
ckpt_manager.thermal_check()
# Entrainement avec AMP
train_metrics = train_epoch(model, train_loader, optimizer, scheduler, device, grad_scaler, thermal_fn=ckpt_manager.batch_thermal_check)
# Validation
val_metrics = evaluate(model, val_loader, device)
epoch_time = time.time() - epoch_start
# Sauvegarder l'historique
for key in history:
prefix = 'train_' if key.startswith('train') else 'val_'
metric_key = key.replace(prefix, '')
metrics = train_metrics if prefix == 'train_' else val_metrics
history[key].append(metrics[metric_key])
# GPU temp pour affichage
gpu_temp = get_gpu_temp()
temp_str = f" | GPU: {gpu_temp}C" if gpu_temp > 0 else ""
# Afficher le progres
print(
f"Epoch {epoch+1:2d}/{EPOCHS} | "
f"Train Loss: {train_metrics['loss']:.4f} "
f"(reg={train_metrics['reg_loss']:.4f}, cls={train_metrics['cls_loss']:.4f}) | "
f"Val Loss: {val_metrics['loss']:.4f} | "
f"Dir Acc: {val_metrics['dir_acc']:.2%} | "
f"Time: {epoch_time:.1f}s{temp_str}"
)
# Checkpoint update
is_best = ckpt_manager.update(
epoch, val_metrics['loss'], history,
model, optimizer, scheduler, grad_scaler
)
if is_best:
print(f" -> Nouveau meilleur modele (val_loss={val_metrics['loss']:.4f})")
# Early stopping -- plancher MIN_EPOCHS : le controle du 25/09 (deux
# tirages de main non modifie) coupe a 9 puis 17 epochs, et une coupe
# precoce fige un modele aux predictions constantes (Correlation nan).
if ckpt_manager.should_stop():
if (epoch + 1) < MIN_EPOCHS:
print(f"\nEarly stopping differe : plancher MIN_EPOCHS={MIN_EPOCHS} non atteint ({epoch+1} epochs)")
else:
print(f"\nEarly stopping a l'epoch {epoch+1} (patience={ckpt_manager.patience})")
break
# Liberer memoire GPU
if torch.cuda.is_available():
torch.cuda.empty_cache()
# Sauvegarder le modele final
ckpt_manager.finalize(model, extra={
'scaler': scaler,
'config': {
'input_dim': n_features,
'd_model': D_MODEL, 'nhead': NHEAD,
'num_layers': NUM_LAYERS,
'dim_feedforward': DIM_FEEDFORWARD,
'dropout': DROPOUT,
'seq_len': SEQ_LEN, 'pred_len': PRED_LEN,
'feature_cols': feature_cols,
'tickers': TICKERS_30,
}
})
print(f"\nEntrainement termine.")
print(f"Meilleure val_loss: {ckpt_manager.best_val_loss:.4f}")
print(f"Checkpoint de reprise: {os.path.basename(checkpoint_path)} (ne pas commit)")
else:
print("\nModele deja entraine et charge. Aucun entrainement necessaire.")