laya-vulnerability-groups / train_laya.py
peter2000's picture
Add training script: RLCD fine-tune of laya on GIZ vulnerability data (17 binary questions per text)
c9604b1 verified
Raw History Blame Contribute Delete
17.7 kB
import os
os.environ.setdefault("USE_TF", "0")
os.environ.setdefault("USE_TORCH", "1")
os.environ.setdefault("TOKENIZERS_PARALLELISM", "false")
os.environ.setdefault("PYTORCH_CUDA_ALLOC_CONF", "expandable_segments:True")
import json
import random
import shutil
import time
import numpy as np
import pandas as pd
import torch
import trackio
def _trackio_safe(fn):
def wrapper(*a, **k):
try:
return fn(*a, **k)
except Exception as e:
print(f"trackio.{fn.__name__} failed (non-fatal): {e}")
return wrapper
trackio.log = _trackio_safe(trackio.log)
trackio.init = _trackio_safe(trackio.init)
trackio.finish = _trackio_safe(trackio.finish)
from huggingface_hub import HfApi, snapshot_download
from safetensors.torch import load_file, save_file
from sklearn.metrics import f1_score
from sklearn.model_selection import train_test_split
from transformers import AutoTokenizer
import laya
from laya.agent import _fix_tokenizer_config
from laya.common import QTYPES, build_model, build_sequence, proper_reward, render_options
SMOKE = os.environ.get("SMOKE", "0") == "1"
MODEL_ID = "convaiinnovations/laya"
OUTPUT_DIR = os.environ.get("LAYA_OUTPUT_DIR", "/data/laya_vulnerability_groups")
PARQUET_URL = "https://huggingface.co/datasets/GIZ/vulnerability_training_data_full/resolve/refs%2Fconvert%2Fparquet/default/train/0000.parquet"
LABELS = [
"Agricultural communities", "Coastal communities", "Ethnic, racial or other minorities",
"Fishery communities", "Informal sector workers", "Members of indigenous and local communities",
"Migrants and displaced persons", "Older persons", "Other", "Persons living in poverty",
"Persons with disabilities", "Persons with pre-existing health conditions",
"Residents of drought-prone regions", "Rural populations", "Sexual minorities (LGBTQI+)",
"Urban populations", "Women and other genders",
]
REPO_ID = "peter2000/laya-vulnerability-groups"
QIDS = [f"g{i}" for i in range(len(LABELS))]
def make_questions():
return {
qid: {
"type": "noul",
"instructions": f"Does this text indicate that {label} are targeted, supported, or affected as a vulnerable group? Answer true or false.",
}
for qid, label in zip(QIDS, LABELS)
}
QUESTIONS = make_questions()
def load_data():
df = pd.read_parquet(PARQUET_URL)
assert len(df) == 475, f"expected 475 rows, got {len(df)}"
Y = df[LABELS].values.astype(np.int64)
nlab = Y.sum(1)
idx_tr, idx_te = train_test_split(
np.arange(len(df)), test_size=0.2, random_state=42, stratify=np.minimum(nlab, 3)
)
return df, Y, np.asarray(idx_tr), np.asarray(idx_te)
def ece(conf, correct, n_bins=15):
conf = np.asarray(conf, dtype=np.float64)
corr = np.asarray(correct, dtype=np.float64)
bins = np.linspace(0.0, 1.0, n_bins + 1)
e = 0.0
for lo, hi in zip(bins[:-1], bins[1:]):
m = (conf > lo) & (conf <= hi)
if m.sum() > 0:
e += m.mean() * abs(corr[m].mean() - conf[m].mean())
return float(e)
def evaluate(Y_true, P_pred, threshold=0.5):
pred = (P_pred >= threshold).astype(int)
per_label = f1_score(Y_true, pred, average=None, zero_division=0)
macro = float(f1_score(Y_true, pred, average="macro", zero_division=0))
micro = float(f1_score(Y_true, pred, average="micro", zero_division=0))
conf = np.where(pred == 1, P_pred, 1.0 - P_pred)
corr = (pred == Y_true).astype(np.float64)
return {
"macro_f1": macro,
"micro_f1": micro,
"ece": ece(conf, corr),
"subset_accuracy": float(((pred == Y_true).all(axis=1)).mean()),
"per_label_f1": {LABELS[i]: round(float(per_label[i]), 4) for i in range(len(LABELS))},
}
def probs_from_answers(res):
return np.array([res["answers"][qid]["noul"] for qid in QIDS], dtype=np.float64)
def eval_agent(agent, texts, Y_true):
t0 = time.time()
P = np.stack([probs_from_answers(agent.predict(t, QUESTIONS)) for t in texts])
metrics = evaluate(Y_true, P)
metrics["eval_seconds"] = round(time.time() - t0, 1)
return metrics
def collate_train_batch(items, pad_id):
n, L = len(items), max(len(it["ids"]) for it in items)
kmax = max(len(it["markers"]) for it in items)
ids = torch.full((n, L), pad_id, dtype=torch.long)
att = torch.zeros((n, L), dtype=torch.long)
mpos = torch.zeros((n, kmax), dtype=torch.long)
mmask = torch.zeros((n, kmax), dtype=torch.bool)
target = torch.zeros((n, kmax), dtype=torch.float32)
for i, it in enumerate(items):
ids[i, : len(it["ids"])] = torch.tensor(it["ids"])
att[i, : len(it["ids"])] = 1
k = len(it["markers"])
mpos[i, :k] = torch.tensor(it["markers"])
mmask[i, :k] = True
target[i, : len(it["target"])] = torch.tensor(it["target"], dtype=torch.float32)
return {
"input_ids": ids, "attention_mask": att, "marker_pos": mpos,
"marker_mask": mmask, "target": target,
"qtype": torch.tensor([it["qtype"] for it in items]),
"label": torch.tensor([it["label"] for it in items]),
}
def fit_one_temp(sel):
if len(sel) < 10:
return 1.0
kmax = max(len(z) for z, _ in sel)
Z = torch.full((len(sel), kmax), -1e4)
T = torch.zeros((len(sel), kmax))
for i, (z, t) in enumerate(sel):
Z[i, :len(z)] = torch.tensor(z)
T[i, :len(t)] = torch.tensor(t, dtype=torch.float32)
log_t = torch.zeros(1, requires_grad=True)
opt = torch.optim.LBFGS([log_t], lr=0.1, max_iter=100)
def closure():
opt.zero_grad()
loss = -(T * torch.log_softmax(Z / log_t.exp(), -1)).sum(-1).mean()
loss.backward()
return loss
opt.step(closure)
return float(torch.clamp(log_t.exp(), 0.1, 10.0).item())
def build_items(texts, yvecs, tok, cfg):
items, dropped = [], 0
k_expected = len(render_options({"t": "noul", "crit": {}}))
for text, yvec in zip(texts, yvecs):
for j in range(len(LABELS)):
p_true = float(yvec[j])
target = [1.0 - p_true, p_true]
q = {"t": "noul", "ins": QUESTIONS[QIDS[j]]["instructions"], "crit": {}}
seq, markers = build_sequence(tok, text, q, cfg["max_len"], cfg["head_max_len"])
if len(markers) != k_expected:
dropped += 1
continue
items.append({
"ids": seq, "markers": markers, "qtype": QTYPES["noul"],
"target": target, "label": j,
})
print(f"built {len(items)} items, dropped {dropped} (marker mismatch), k_expected={k_expected}")
return items
def train(items, model, tok, cfg, device):
EPOCHS = 1 if SMOKE else 4
MICRO_BATCH = 8
GRAD_ACCUM = 4
GROUP_SIZE = 4
LR_ENCODER = 2.5e-5
LR_HEAD = 1.0e-4
SIGMA_START = 0.4
SIGMA_END = 0.1
enc_params = [p for n, p in model.named_parameters() if "encoder." in n]
head_params = [p for n, p in model.named_parameters() if "encoder." not in n]
optimizer = torch.optim.AdamW([
{"params": enc_params, "lr": LR_ENCODER},
{"params": head_params, "lr": LR_HEAD},
], weight_decay=0.01)
total_updates = max(1, (len(items) // (MICRO_BATCH * GRAD_ACCUM)) * EPOCHS)
scheduler = torch.optim.lr_scheduler.CosineAnnealingLR(optimizer, T_max=total_updates, eta_min=1e-6)
scaler = torch.amp.GradScaler("cuda", enabled=True)
pad_id = tok.pad_token_id
t0 = time.time()
for epoch in range(EPOCHS):
random.seed(42 + epoch)
random.shuffle(items)
epoch_loss, n_batches = 0.0, 0
optimizer.zero_grad(set_to_none=True)
accum_step = 0
progress = epoch / max(1, EPOCHS - 1)
sigma = SIGMA_START + (SIGMA_END - SIGMA_START) * progress
for b_idx in range(0, len(items), MICRO_BATCH):
chunk = items[b_idx:b_idx + MICRO_BATCH]
if not chunk:
continue
batch = collate_train_batch(chunk, pad_id)
with torch.autocast("cuda", dtype=torch.float16):
logits, act = model(
batch["input_ids"].to(device),
batch["attention_mask"].to(device),
batch["marker_pos"].to(device),
batch["marker_mask"].to(device),
batch["qtype"].to(device),
)
logits = logits.float()
mask = batch["marker_mask"].to(device)
k = mask.sum(-1, keepdim=True).float()
target = batch["target"].to(device)
eps = torch.randn((GROUP_SIZE,) + logits.shape, device=device) * sigma * mask
eps = (eps - eps.sum(-1, keepdim=True) / k) * mask
z = logits.detach().unsqueeze(0) + eps
q = torch.softmax(z.masked_fill(~mask, -1e4), -1)
with torch.no_grad():
r = proper_reward(q, target.unsqueeze(0), batch["qtype"].to(device), mask, w_sph=0.75, w_rps=1.0)
adv = r - r.mean(0, keepdim=True)
adv = adv / (adv.std() + 1e-6)
logp = -(((z - logits.unsqueeze(0)) ** 2) * mask).sum(-1) / (2 * sigma ** 2)
loss_rl = -(adv * logp).mean()
loss_ce = -(target * torch.log_softmax(logits.masked_fill(~mask, -1e4), -1)).sum(-1).mean()
loss = (loss_rl + 1.0 * loss_ce) / GRAD_ACCUM + 0.0 * act.sum()
scaler.scale(loss).backward()
accum_step += 1
if accum_step % GRAD_ACCUM == 0 or (b_idx + MICRO_BATCH) >= len(items):
scaler.unscale_(optimizer)
torch.nn.utils.clip_grad_norm_(model.parameters(), 1.0)
scaler.step(optimizer)
scaler.update()
scheduler.step()
optimizer.zero_grad(set_to_none=True)
epoch_loss += loss.item() * GRAD_ACCUM
n_batches += 1
if n_batches % 50 == 0:
print(f" epoch {epoch+1}/{EPOCHS} step {n_batches} loss {loss.item()*GRAD_ACCUM:.4f} reward {r.mean().item():.3f} lr {scheduler.get_last_lr()[0]:.2e}")
trackio.log({"laya_loss": loss.item() * GRAD_ACCUM, "laya_reward": r.mean().item()}, step=epoch * 100000 + n_batches)
print(f"=== epoch {epoch+1}/{EPOCHS} done in {time.time()-t0:.0f}s avg_loss {epoch_loss/max(1,n_batches):.4f} ===")
trackio.log({"laya_epoch_loss": epoch_loss / max(1, n_batches)}, step=epoch + 1)
ckpt_dir = os.path.join(OUTPUT_DIR, "checkpoint_latest")
os.makedirs(ckpt_dir, exist_ok=True)
ckpt_sd = {k: v.half().contiguous().cpu() for k, v in model.state_dict().items()}
save_file(ckpt_sd, os.path.join(ckpt_dir, "model.safetensors"))
model.encoder.config.save_pretrained(os.path.join(ckpt_dir, "encoder"))
tok.save_pretrained(os.path.join(ckpt_dir, "tokenizer"))
with open(os.path.join(ckpt_dir, "checkpoint_meta.json"), "w") as f:
json.dump({"epoch": epoch + 1, "total_epochs": EPOCHS, "avg_loss": epoch_loss / max(1, n_batches)}, f, indent=2)
del optimizer, scaler, scheduler
torch.cuda.empty_cache()
return model
def fit_and_apply_temperature(items, model, cfg, device, pad_id):
model.eval()
calib_items = items[::15][:400]
calib_preds = []
with torch.no_grad():
for c_idx in range(0, len(calib_items), 16):
c_chunk = calib_items[c_idx:c_idx + 16]
cb = collate_train_batch(c_chunk, pad_id)
with torch.autocast("cuda", dtype=torch.float16):
l_sub, _ = model(
cb["input_ids"].to(device),
cb["attention_mask"].to(device),
cb["marker_pos"].to(device),
cb["marker_mask"].to(device),
cb["qtype"].to(device),
)
l_np = l_sub.float().cpu().numpy()
for rr, it in enumerate(c_chunk):
k = len(it["markers"])
calib_preds.append((it["qtype"], l_np[rr, :k], it["target"]))
fitted = [1.2, 1.2, 1.2]
for qt in range(3):
sel = [(z, t) for q_type, z, t in calib_preds if q_type == qt]
if sel:
fitted[qt] = fit_one_temp(sel)
print("fitted calibration temperatures (choice, score, noul):", [round(t, 3) for t in fitted])
cfg["temperature"] = fitted
if isinstance(cfg.get("temperature_by_options"), dict):
cfg["temperature_by_options"]["noul:2"] = fitted[QTYPES["noul"]]
return fitted
def save_checkpoint(model, tok, cfg):
os.makedirs(OUTPUT_DIR, exist_ok=True)
sd = {k: v.half().contiguous().cpu() for k, v in model.state_dict().items()}
save_file(sd, os.path.join(OUTPUT_DIR, "model.safetensors"))
model.encoder.config.save_pretrained(os.path.join(OUTPUT_DIR, "encoder"))
tok.save_pretrained(os.path.join(OUTPUT_DIR, "tokenizer"))
cfg["fine_tuned"] = True
cfg["model_name"] = "laya-vulnerability-groups"
with open(os.path.join(OUTPUT_DIR, "rl_agent_config.json"), "w") as f:
json.dump(cfg, f, indent=2)
ckpt_dir = os.path.join(OUTPUT_DIR, "checkpoint_latest")
if os.path.isdir(ckpt_dir):
shutil.rmtree(ckpt_dir)
print(f"checkpoint saved to {OUTPUT_DIR}")
def write_readme(metrics):
readme = f"""---
license: apache-2.0
library_name: transformers
tags: [laya, text-classification, multi-label, climate, vulnerability, rlcd]
pipeline_tag: text-classification
---
# Laya fine-tuned for climate-vulnerability group detection (multi-label)
[convaiinnovations/laya](https://huggingface.co/convaiinnovations/laya) (421M, ModernBERT-large backbone)
fine-tuned with the official RLCD recipe on [GIZ/vulnerability_training_data_full](https://huggingface.co/datasets/GIZ/vulnerability_training_data_full)
(380 train rows x 17 binary vulnerability-group questions; 36 all-negative rows included as negatives).
Each vulnerability group is asked as one binary (`noul`) typed question; all 17 are answered in a single forward pass.
Evaluate with `laya.load("peter2000/laya-vulnerability-groups")` and `agent.predict(state, questions)`.
## Test-set metrics (held-out {95 if not SMOKE else 5} rows, threshold 0.5)
| metric | value |
|---|---|
| macro-F1 | {metrics['macro_f1']:.4f} |
| micro-F1 | {metrics['micro_f1']:.4f} |
| ECE | {metrics['ece']:.4f} |
| subset accuracy | {metrics['subset_accuracy']:.4f} |
Per-label F1:
| label | F1 |
|---|---|
"""
for label, f1 in metrics["per_label_f1"].items():
readme += f"| {label} | {f1:.4f} |\n"
with open(os.path.join(OUTPUT_DIR, "README.md"), "w") as f:
f.write(readme)
def main():
df, Y, idx_tr, idx_te = load_data()
texts = df["text"].tolist()
X_tr = [texts[i] for i in idx_tr]
X_te = [texts[i] for i in idx_te]
Y_tr, Y_te = Y[idx_tr], Y[idx_te]
if SMOKE:
X_te, Y_te = X_te[:5], Y_te[:5]
print(f"train={len(X_tr)} test={len(X_te)} labels={len(LABELS)}")
trackio.init(project="vulnerability-multilabel-classifier", space_id="peter2000/vulnerability-multilabel-classifier-trackio")
model_dir = snapshot_download(MODEL_ID, ignore_patterns=["multilingual/*", "typed-decisions/*", "assets/*", "eval/*", "*.py"])
_fix_tokenizer_config(model_dir)
tok = AutoTokenizer.from_pretrained(os.path.join(model_dir, "tokenizer"))
with open(os.path.join(model_dir, "rl_agent_config.json")) as f:
cfg = json.load(f)
cfg["gradient_checkpointing"] = True
device = "cuda"
print("--- zero-shot baseline eval ---")
agent0 = laya.load(model_dir, device=device)
m0 = eval_agent(agent0, X_te, Y_te)
print("zero-shot:", json.dumps(m0, indent=2))
trackio.log({"laya_zeroshot_macro_f1": m0["macro_f1"], "laya_zeroshot_micro_f1": m0["micro_f1"], "laya_zeroshot_ece": m0["ece"]}, step=0)
del agent0
torch.cuda.empty_cache()
items = build_items(X_tr, Y_tr, tok, cfg)
if SMOKE:
items = items[:64]
print("building model...")
model = build_model(cfg, encoder_dir=os.path.join(model_dir, "encoder"))
model.load_state_dict(load_file(os.path.join(model_dir, "model.safetensors")), strict=True)
model.encoder.gradient_checkpointing_enable(gradient_checkpointing_kwargs={"use_reentrant": False})
model.head_checkpointing = True
model.to(device)
model.train()
train(items, model, tok, cfg, device)
fitted = fit_and_apply_temperature(items, model, cfg, device, tok.pad_token_id)
print("fitted temperatures:", [round(t, 3) for t in fitted])
save_checkpoint(model, tok, cfg)
print("--- fine-tuned eval ---")
agent_ft = laya.load(OUTPUT_DIR, device=device)
m1 = eval_agent(agent_ft, X_te, Y_te)
print("fine-tuned:", json.dumps(m1, indent=2))
trackio.log({"laya_macro_f1": m1["macro_f1"], "laya_micro_f1": m1["micro_f1"], "laya_ece": m1["ece"]}, step=10)
trackio.log({f"laya_f1/{k}": v for k, v in m1["per_label_f1"].items()}, step=10)
write_readme(m1)
api = HfApi(token=os.environ.get("HF_TOKEN"))
api.upload_folder(
folder_path=OUTPUT_DIR,
repo_id=REPO_ID,
repo_type="model",
commit_message=f"Laya fine-tuned on GIZ vulnerability data: macro-F1 {m1['macro_f1']:.3f} (zero-shot {m0['macro_f1']:.3f})",
)
api.upload_file(
path_or_fileobj=json.dumps({"zero_shot": m0, "fine_tuned": m1}, indent=2).encode(),
path_in_repo="metrics.json",
repo_id=REPO_ID,
repo_type="model",
commit_message="Add evaluation metrics",
)
trackio.finish()
print("DONE")
if __name__ == "__main__":
t0 = time.time()
main()
print(f"elapsed {time.time()-t0:.0f}s")