Files

316 lines
17 KiB
Python

"""Evaluate the Q1 V2 B0-B4 branches with leakage-safe video-group folds."""
from __future__ import annotations
import argparse
import csv
import hashlib
import json
import math
from pathlib import Path
from typing import Any
import numpy as np
from sklearn.linear_model import LogisticRegression, Ridge
from sklearn.metrics import accuracy_score, f1_score, mean_absolute_error, mean_squared_error
from sklearn.model_selection import GroupKFold
from . import compare_models as cm
from . import q1_io
Q1 = Path(__file__).resolve().parent
OUT = Q1.parent / "output" / "q1" / "model_comparison_v2"
MODALITIES = ("text", "audio", "vision", "speech")
METHODS = ("B0", "B1", "B2", "B3", "B4")
def digest(path: Path) -> str:
h = hashlib.sha256()
with path.open("rb") as stream:
for block in iter(lambda: stream.read(1024 * 1024), b""):
h.update(block)
return h.hexdigest()
def edges_from_bounds(bounds: np.ndarray) -> np.ndarray:
return np.concatenate((bounds[:1, 0], bounds[:, 1])).astype(np.float32)
def as_view(views: dict[str, dict[str, np.ndarray]], bounds: np.ndarray) -> cm.View:
return cm.View(
features={name: views[name]["x"].astype(np.float32) for name in MODALITIES},
observed={name: views[name]["mask"].astype(bool) for name in MODALITIES},
coverage={name: views[name]["coverage"].astype(np.float32) for name in MODALITIES},
edges=edges_from_bounds(bounds),
)
def load_rows(records: list[dict[str, Any]]) -> tuple[list[dict[str, Any]], list[cm.View], list[dict[str, np.ndarray]], list[cm.View], list[cm.View]]:
rows, views, dynamics = [], [], []
feature_ids = {json.loads(line)["sample_id"] for line in (q1_io.FEATURE_DIR / "manifest_q1.jsonl").read_text(encoding="utf-8").splitlines()}
record_ids = [record["sample_id"] for record in records]
if len(set(record_ids)) != len(record_ids) or set(record_ids) != feature_ids:
raise ValueError("Q1 labels and V2 feature manifest keys are not a one-to-one match")
for record in records:
sample = q1_io.load_sample(record["sample_id"])
base = {name: q1_io.get_view(sample, name, "sec") for name in ("text", "audio", "vision")}
base["speech"] = q1_io.get_view(sample, "speech", "sec")
vision_b1 = q1_io.get_view(sample, "vision", "sec", geometry_aware=True)
b1 = {**base, "vision": vision_b1}
posterior = {**b1, "text": q1_io.get_view(sample, "text", "posterior")}
rows.append({**record, "duration_s": float(sample["_meta"]["duration_s"])})
# Every model branch shares the same physical primary time bounds.
bounds = sample["views_sec_time_bounds_s"].astype(np.float32)
views.append(as_view(base, bounds))
views.append(as_view(b1, bounds))
views.append(as_view(posterior, bounds))
with np.load(sample["_path"], allow_pickle=False) as artifact:
dynamics.append({"x": artifact["views_dynamics_x"].astype(np.float32), "mask": artifact["views_dynamics_mask"].astype(bool), "bounds": artifact["views_dynamics_time_bounds_s"].astype(np.float32)})
# Flatten rows as [B0,B1,B2] per sample into three separate lists.
n = len(records)
return rows, [views[3*i] for i in range(n)], dynamics, [views[3*i+1] for i in range(n)], [views[3*i+2] for i in range(n)]
def fit_scaler(views: list[cm.View], indices: np.ndarray, mode: str) -> dict[str, tuple[np.ndarray, np.ndarray]]:
scalers = {}
for name in MODALITIES:
d = views[int(indices[0])].features[name].shape[1]
center = np.zeros(d, np.float64)
scale = np.ones(d, np.float64)
for j in range(d):
chunks = [views[int(i)].features[name][:, j][views[int(i)].observed[name][:, j]] for i in indices]
chunks = [x for x in chunks if len(x)]
if not chunks:
continue
values = np.concatenate(chunks).astype(np.float64)
if mode == "robust":
center[j] = np.median(values)
mad = 1.4826 * np.median(np.abs(values - center[j]))
scale[j] = mad if mad > 1e-8 else (np.std(values) if np.std(values) > 1e-8 else 1.0)
else:
center[j] = np.mean(values)
scale[j] = np.std(values) if np.std(values) > 1e-8 else 1.0
scalers[name] = (center.astype(np.float32), scale.astype(np.float32))
return scalers
def scale_view(view: cm.View, scalers: dict[str, tuple[np.ndarray, np.ndarray]], robust: bool) -> cm.View:
result = {}
for name in MODALITIES:
center, scale = scalers[name]
values = (view.features[name] - center) / scale
if robust:
values = np.clip(values, -8.0, 8.0)
result[name] = np.where(view.observed[name], values, 0.0).astype(np.float32)
return cm.View(result, view.observed, view.coverage, view.edges)
def pool_progress(view: cm.View, include_speech: bool) -> np.ndarray:
duration = float(view.edges[-1])
progress = np.linspace(0.0, duration, 6)
centers = (view.edges[:-1] + view.edges[1:]) / 2
lengths = np.diff(view.edges)
pooled = []
names = ("text", "audio", "vision") + (("speech",) if include_speech else ())
for name in names:
values = view.features[name]
mask = view.observed[name]
coverage = view.coverage[name]
segments = []
for left, right in zip(progress[:-1], progress[1:]):
overlap = np.maximum(0.0, np.minimum(view.edges[1:], right) - np.maximum(view.edges[:-1], left))
weight = overlap[:, None] * coverage * mask
den = weight.sum(axis=0)
out = np.zeros(values.shape[1], np.float32)
good = den > 0
if good.any():
out[good] = (values * weight).sum(axis=0)[good] / den[good]
segments.append(out)
pooled.append(np.stack(segments).reshape(-1))
return np.concatenate(pooled).astype(np.float32)
def pool_dynamics(item: dict[str, np.ndarray], duration: float) -> np.ndarray:
bounds = item["bounds"]
values = item["x"]
valid = item["mask"]
centers = bounds[:, 0].mean(axis=1)
grid_length = bounds[:, 0, 1] - bounds[:, 0, 0]
progress = np.linspace(0.0, duration, 6)
result = []
for left, right in zip(progress[:-1], progress[1:]):
sel = (centers >= left) & (centers < right)
for branch in range(2):
ok = sel & valid[:, branch]
weights = grid_length[ok]
if ok.any() and weights.sum() > 0:
result.extend(np.average(values[ok, branch], axis=0, weights=weights).tolist())
result.append(float(min(1.0, weights.sum() / max(right - left, 1e-8))))
else:
result.extend([0.0] * 6)
result.append(0.0)
return np.asarray(result, np.float32)
def fit_dynamic_scaler(values: np.ndarray, indices: np.ndarray) -> tuple[np.ndarray, np.ndarray, np.ndarray]:
center = np.zeros(values.shape[1], np.float32)
scale = np.ones(values.shape[1], np.float32)
is_signature = np.ones(values.shape[1], bool)
# In each 14-column block, the final column is a measured coverage fraction.
is_signature[np.arange(6, values.shape[1], 7)] = False
for j in np.flatnonzero(is_signature):
coverage_j = (j // 7) * 7 + 6
observed = values[indices, coverage_j] > 0
vals = values[indices[observed], j]
if not len(vals):
continue
center[j] = np.median(vals)
mad = 1.4826 * np.median(np.abs(vals - center[j]))
scale[j] = mad if mad > 1e-8 else (np.std(vals) if np.std(vals) > 1e-8 else 1.0)
return center, scale, is_signature
def scale_dynamic(values: np.ndarray, params: tuple[np.ndarray, np.ndarray, np.ndarray]) -> np.ndarray:
center, scale, signature = params
result = values.copy()
result[:, signature] = np.clip((result[:, signature] - center[signature]) / scale[signature], -8.0, 8.0)
# Preserve zeros in missing signature blocks.
for block in range(10):
coverage_col = block * 7 + 6
result[result[:, coverage_col] <= 0, block * 7:block * 7 + 6] = 0.0
return result.astype(np.float32)
def metrics(y_class: np.ndarray, p_class: np.ndarray, y_value: np.ndarray, p_value: np.ndarray) -> dict[str, float]:
pearson = float(np.corrcoef(y_value, p_value)[0, 1]) if np.std(p_value) > 1e-12 else math.nan
return {
"accuracy": float(accuracy_score(y_class, p_class)),
"macro_f1": float(f1_score(y_class, p_class, labels=[0, 1, 2], average="macro", zero_division=0)),
"mae": float(mean_absolute_error(y_value, p_value)),
"rmse": float(np.sqrt(mean_squared_error(y_value, p_value))),
"pearson": pearson,
}
def write_csv(path: Path, rows: list[dict[str, Any]]) -> None:
if not rows:
return
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("w", encoding="utf-8-sig", newline="") as stream:
writer = csv.DictWriter(stream, fieldnames=list(rows[0]))
writer.writeheader()
writer.writerows(rows)
def main() -> None:
global OUT
parser = argparse.ArgumentParser()
parser.add_argument("--bootstrap-repeats", type=int, default=2000)
parser.add_argument("--data-dir", type=Path, default=cm.DATA_DIR)
parser.add_argument("--feature-dir", type=Path, default=q1_io.FEATURE_DIR)
parser.add_argument("--output-dir", type=Path, default=OUT)
args = parser.parse_args()
cm.DATA_DIR = args.data_dir.resolve()
cm.LABEL_FILE = cm.DATA_DIR / "label-100.xlsx"
q1_io.FEATURE_DIR = args.feature_dir.resolve()
OUT = args.output_dir.resolve()
records = cm._read_labels(cm.LABEL_FILE)
rows, b0, dynamics, b1, b2 = load_rows(records)
y_class = np.asarray([r["polarity"] for r in rows], np.int64)
y_value = np.asarray([r["sentiment"] for r in rows], np.float64)
groups = np.asarray([r["video_id"] for r in rows], str)
durations = np.asarray([r["duration_s"] for r in rows], np.float64)
dynamic_values = np.stack([pool_dynamics(d, durations[i]) for i, d in enumerate(dynamics)])
variants = {"B0": b0, "B1": b1, "B2": b2, "B3": b1, "B4": b1}
splitter = GroupKFold(n_splits=cm.N_FOLDS)
folds = list(splitter.split(np.zeros(len(rows)), y_class, groups))
class_pred = {m: np.full(len(rows), -1, np.int64) for m in METHODS}
value_pred = {m: np.full(len(rows), np.nan, np.float64) for m in METHODS}
fold_rows = []
split_rows = []
for fold, (train_idx, valid_idx) in enumerate(folds, start=1):
if set(groups[train_idx]) & set(groups[valid_idx]):
raise ValueError("video group leakage in Q1 folds")
for i in valid_idx:
split_rows.append({"sample_id": rows[i]["sample_id"], "video_id": groups[i], "fold": fold, "split": "valid_oof"})
for method in METHODS:
source = variants[method]
robust = method != "B0"
scaler = fit_scaler(source, train_idx, "robust" if robust else "standard")
scaled = [scale_view(v, scaler, robust) for v in source]
include_speech = method == "B4"
x_train = np.stack([pool_progress(scaled[i], include_speech) for i in train_idx])
x_valid = np.stack([pool_progress(scaled[i], include_speech) for i in valid_idx])
if method == "B3":
dyn_scaler = fit_dynamic_scaler(dynamic_values, train_idx)
dyn_scaled = scale_dynamic(dynamic_values, dyn_scaler)
x_train = np.column_stack((x_train, dyn_scaled[train_idx]))
x_valid = np.column_stack((x_valid, dyn_scaled[valid_idx]))
clf = LogisticRegression(C=0.05, max_iter=2500, solver="lbfgs", random_state=cm.SEED)
clf.fit(x_train, y_class[train_idx])
p_cls = clf.predict(x_valid)
reg = Ridge(alpha=25.0, solver="lsqr")
reg.fit(x_train, y_value[train_idx])
p_val = np.clip(reg.predict(x_valid), -3.0, 3.0)
class_pred[method][valid_idx] = p_cls
value_pred[method][valid_idx] = p_val
met = metrics(y_class[valid_idx], p_cls, y_value[valid_idx], p_val)
fold_rows.append({"method": method, "fold": fold, "train_samples": len(train_idx), "valid_samples": len(valid_idx), "train_video_groups": len(set(groups[train_idx])), "valid_video_groups": len(set(groups[valid_idx])), "feature_dimension": int(x_train.shape[1]), **met})
print(f"fold {fold}/{cm.N_FOLDS} complete", flush=True)
summary = []
prediction_rows = []
for method in METHODS:
overall = metrics(y_class, class_pred[method], y_value, value_pred[method])
fold = [r for r in fold_rows if r["method"] == method]
row = {"method": method, "sample_count": len(rows), "video_group_count": len(set(groups)), **{f"oof_{k}": v for k, v in overall.items()}}
for key in ("accuracy", "macro_f1", "mae", "rmse", "pearson"):
vals = np.asarray([r[key] for r in fold], np.float64)
row[f"fold_{key}_mean"] = float(np.nanmean(vals))
row[f"fold_{key}_sd"] = float(np.nanstd(vals, ddof=1))
summary.append(row)
for i, record in enumerate(rows):
line = {"sample_id": record["sample_id"], "video_id": record["video_id"], "clip_id": record["clip_id"], "true_polarity": int(y_class[i]), "true_sentiment": float(y_value[i])}
for method in METHODS:
line[f"{method}_predicted_polarity"] = int(class_pred[method][i])
line[f"{method}_predicted_sentiment"] = float(value_pred[method][i])
prediction_rows.append(line)
rng = np.random.default_rng(cm.SEED + 880)
unique_groups = np.unique(groups)
group_ix = {g: np.flatnonzero(groups == g) for g in unique_groups}
deltas = []
reference = "B1"
for method in METHODS:
if method == reference:
continue
samples = {metric: [] for metric in ("accuracy", "macro_f1", "mae", "rmse", "pearson")}
for _ in range(args.bootstrap_repeats):
selected_groups = rng.choice(unique_groups, size=len(unique_groups), replace=True)
ix = np.concatenate([group_ix[g] for g in selected_groups])
base = metrics(y_class[ix], class_pred[reference][ix], y_value[ix], value_pred[reference][ix])
alt = metrics(y_class[ix], class_pred[method][ix], y_value[ix], value_pred[method][ix])
for key in samples:
samples[key].append(alt[key] - base[key])
for key, values in samples.items():
vals = np.asarray(values, np.float64)
deltas.append({"comparison": f"{method}-{reference}", "metric": key, "reference": reference, "estimate_delta": float(np.nanmedian(vals)), "ci95_low": float(np.nanpercentile(vals, 2.5)), "ci95_high": float(np.nanpercentile(vals, 97.5)), "bootstrap_repeats": args.bootstrap_repeats, "unit": "video_id"})
OUT.mkdir(parents=True, exist_ok=True)
write_csv(OUT / "comparison_summary.csv", summary)
write_csv(OUT / "fold_metrics.csv", fold_rows)
write_csv(OUT / "oof_predictions.csv", prediction_rows)
write_csv(OUT / "split_assignments.csv", sorted(split_rows, key=lambda r: r["sample_id"]))
write_csv(OUT / "group_bootstrap_deltas.csv", deltas)
manifest = {"version": "E题V2 Q1 B0-B4", "samples": len(rows), "video_groups": int(len(unique_groups)), "folds": cm.N_FOLDS, "grouping": "video_id GroupKFold; train/valid group disjoint", "scalers": "B0 train-fold mean/std; B1-B4 train-fold median/MAD; no held-fold fitting", "classifier": "LogisticRegression C=0.05", "regressor": "Ridge alpha=25", "branches": {"B0": "quality-weighted physical projection", "B1": "B0 plus train-fold robust scale and SO(3)/unit-gaze aggregation", "B2": "B1 plus CTC forward-backward text occupancy", "B3": "B1 plus native-row level-2 dynamics", "B4": "B1 plus frozen Wav2Vec2 auxiliary representation"}, "B5": "query H_time audit only; not predictive score", "B6": "content probe disabled", "feature_manifest_sha256": digest(q1_io.FEATURE_DIR / "feature_manifest.json"), "source_run_manifest": "features_v2/manifest_q1.jsonl", "bootstrap_repeats": args.bootstrap_repeats}
(OUT / "run_manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2), encoding="utf-8")
report = ["# Q1 V2 B0-B4 对照", "", "本报告使用 `features_v2/` 的 100 个样本,以 `video_id` 做五折 GroupKFold。B0–B4 的缩放器只在各折训练部分拟合。分数是情感信息探针,不是真实对齐边界准确率。", "", "| 分支 | OOF Accuracy | OOF Macro-F1 | OOF MAE | OOF RMSE | OOF Pearson |", "|---|---:|---:|---:|---:|---:|"]
for row in summary:
report.append(f"| {row['method']} | {row['oof_accuracy']:.4f} | {row['oof_macro_f1']:.4f} | {row['oof_mae']:.4f} | {row['oof_rmse']:.4f} | {row['oof_pearson']:.4f} |")
report += ["", "B5 通过词查询 CSR 矩阵进行物理来源核查;没有人工边界真值,因此未报告边界误差或后验校准率。B6 内容探针未启用。视觉 17 维来自 OpenFace 2.2.0 AU 输出,不是 MediaPipe 代理。"]
(OUT / "report.md").write_text("\n".join(report) + "\n", encoding="utf-8")
print("Q1 V2 comparison:", [(r["method"], round(r["oof_accuracy"], 4), round(r["oof_macro_f1"], 4), round(r["oof_mae"], 4)) for r in summary], flush=True)
if __name__ == "__main__":
main()