from pathlib import Path import sys,os,torch import torch.distributed as dist from torch.nn.parallel import DistributedDataParallel as DDP ROOT=Path(__file__).resolve().parents[1];sys.path.insert(0,str(ROOT)) from model.global_flood_lstm import * c=load_config(ROOT);rank=int(os.getenv("RANK",0));world=int(os.getenv("WORLD_SIZE",1));distributed=world>1 if distributed:dist.init_process_group("gloo") torch.manual_seed(c["seed"]);torch.set_num_threads(2);base=GlobalFloodLSTM(**c["model"]);m=DDP(base) if distributed else base;opt=torch.optim.Adam(m.parameters(),lr=c["train"]["learning_rate"]);ids=list(range(rank,c["data"]["basins"],world));losses=[] for _ in range(c["train"]["epochs"]): for i in ids:h,f,s,y=synthetic_sample(i);loc,scale,tau=m(h[None],f[None],s[None]);loss=asymmetric_laplace_nll(y[None],loc,scale,tau);opt.zero_grad();loss.backward();opt.step();losses.append(float(loss)) v=torch.tensor([sum(losses),len(losses)],dtype=torch.float64) if distributed:dist.all_reduce(v) p=ROOT/c["paths"]["checkpoint"] if rank==0:p.parent.mkdir(parents=True,exist_ok=True);torch.save({"model":base.state_dict(),"model_config":c["model"]},p);write_json(ROOT/c["paths"]["training_metrics"],{"negative_log_likelihood":float(v[0]/v[1]),"world_size":world});print(p) if distributed:dist.destroy_process_group()