240 lines
13 KiB
Python
240 lines
13 KiB
Python
from __future__ import annotations
|
|
|
|
import argparse
|
|
import csv
|
|
import hashlib
|
|
import json
|
|
import pickle
|
|
import shutil
|
|
import sys
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
import numpy as np
|
|
|
|
ROOT = Path(__file__).resolve().parents[2]
|
|
Q1_DIR = Path(__file__).resolve().parent
|
|
sys.path.insert(0, str(Q1_DIR))
|
|
import compare_models as cm # noqa: E402
|
|
|
|
CACHE_SCHEMA = "q1-b0b4-v3"
|
|
EXPORT_BINS = 50
|
|
MODALITIES = ("text", "audio", "vision")
|
|
|
|
|
|
def sha256(path: Path) -> str:
|
|
digest = hashlib.sha256()
|
|
with path.open("rb") as stream:
|
|
for chunk in iter(lambda: stream.read(1024 * 1024), b""):
|
|
digest.update(chunk)
|
|
return digest.hexdigest()
|
|
|
|
|
|
def defaults() -> tuple[Path, Path, Path]:
|
|
candidates = (ROOT / "math" / "Q1" / "cache" / "native", Q1_DIR / "cache" / "native")
|
|
cache = next((path for path in candidates if path.exists()), candidates[0])
|
|
labels = ROOT / "E题数据" / "附件1-数据集原始多模态样本" / "MOSEI数据集部分原始视频-100条" / "label-100.xlsx"
|
|
return cache, labels, cache.parent.parent / "results" / "model_comparison" / "run_manifest.json"
|
|
|
|
|
|
def pool_to_edges(view: cm.View, modality: str, target_edges: np.ndarray) -> tuple[np.ndarray, np.ndarray, np.ndarray]:
|
|
values = np.asarray(view.features[modality], dtype=np.float64)
|
|
observed = np.asarray(view.observed[modality], dtype=bool)
|
|
source_coverage = np.asarray(view.coverage[modality], dtype=np.float64)
|
|
dim = values.shape[1]
|
|
count = len(target_edges) - 1
|
|
output = np.zeros((count, dim), dtype=np.float32)
|
|
mask = np.zeros((count, dim), dtype=bool)
|
|
coverage = np.zeros((count, dim), dtype=np.float32)
|
|
source_left, source_right = view.edges[:-1], view.edges[1:]
|
|
for index, (left, right) in enumerate(zip(target_edges[:-1], target_edges[1:])):
|
|
overlap = np.maximum(0.0, np.minimum(source_right, right) - np.maximum(source_left, left))
|
|
weights = overlap[:, None] * source_coverage * observed
|
|
denominator = weights.sum(axis=0)
|
|
good = denominator > 0
|
|
if good.any():
|
|
output[index, good] = (values * weights).sum(axis=0)[good] / denominator[good]
|
|
mask[index, good] = True
|
|
coverage[index, good] = np.minimum(1.0, denominator[good] / max(right - left, 1e-8))
|
|
return output, mask, coverage
|
|
|
|
|
|
def main() -> None:
|
|
cache_default, labels_default, manifest_default = defaults()
|
|
parser = argparse.ArgumentParser(description="Export Q1 multimodal features for model training.")
|
|
parser.add_argument("--cache-dir", type=Path, default=cache_default)
|
|
parser.add_argument("--labels", type=Path, default=labels_default)
|
|
parser.add_argument("--run-manifest", type=Path, default=manifest_default)
|
|
parser.add_argument("--output-dir", type=Path, default=Q1_DIR / "features")
|
|
args = parser.parse_args()
|
|
|
|
for path in (args.cache_dir, args.labels, args.run_manifest):
|
|
if not path.exists():
|
|
raise FileNotFoundError(f"Required input does not exist: {path}")
|
|
run_manifest = json.loads(args.run_manifest.read_text(encoding="utf-8"))
|
|
inputs = run_manifest["inputs"]
|
|
if run_manifest.get("cache_schema", CACHE_SCHEMA) != CACHE_SCHEMA:
|
|
raise ValueError("Run manifest cache_schema does not match the Q1 exporter.")
|
|
label_hash = sha256(args.labels)
|
|
if inputs.get("label_file_sha256") and inputs["label_file_sha256"] != label_hash:
|
|
raise ValueError("Label workbook hash differs from the run manifest.")
|
|
records = cm._read_labels(args.labels)
|
|
source_hashes = inputs["video_sha256_by_sample"]
|
|
if len(records) != 100 or int(inputs.get("sample_count", -1)) != len(records):
|
|
raise ValueError("Expected the same 100 samples used by the Q1 comparison run.")
|
|
|
|
keys = (
|
|
"text", "text_mask", "text_posterior", "text_posterior_mask", "audio", "audio_mask", "vision", "vision_mask",
|
|
"text_coverage", "text_posterior_activity", "audio_coverage", "vision_coverage", "time_bounds_s", "time_s", "progress",
|
|
"duration_s", "valid_lengths", "id", "sample_id", "video_id", "clip_id", "raw_text", "classification_labels",
|
|
"regression_labels", "annotations", "label_consistent", "source_video_sha256",
|
|
)
|
|
rows: dict[str, list[Any]] = {key: [] for key in keys}
|
|
native: dict[str, list[np.ndarray]] = {key: [] for key in (
|
|
"text", "text_mask", "audio", "audio_mask", "vision", "vision_mask", "text_coverage", "audio_coverage", "vision_coverage", "time_bounds_s",
|
|
)}
|
|
offsets = [0]
|
|
native_ids: list[np.ndarray] = []
|
|
durations: list[float] = []
|
|
|
|
for sample_index, record in enumerate(records):
|
|
sample_id = record["sample_id"]
|
|
source_hash = source_hashes[sample_id]
|
|
cache_path = args.cache_dir / f"{cm._safe_name(record['video_id'])}__{cm._safe_name(record['clip_id'])}.npz"
|
|
sample = cm._load_cache(cache_path, record, source_hash, CACHE_SCHEMA)
|
|
_, view_b1, view_b2 = cm._make_views(sample)
|
|
target_edges = np.linspace(0.0, sample.duration_s, EXPORT_BINS + 1, dtype=np.float64)
|
|
pooled = {name: pool_to_edges(view_b1, name, target_edges) for name in MODALITIES}
|
|
posterior = pool_to_edges(view_b2, "text", target_edges)
|
|
|
|
step_count = len(view_b1.edges) - 1
|
|
native["time_bounds_s"].append(np.column_stack((view_b1.edges[:-1], view_b1.edges[1:])).astype(np.float32))
|
|
native_ids.append(np.full(step_count, sample_index, dtype=np.int16))
|
|
offsets.append(offsets[-1] + step_count)
|
|
durations.append(float(sample.duration_s))
|
|
for name in MODALITIES:
|
|
native[name].append(np.asarray(view_b1.features[name], dtype=np.float16))
|
|
native[f"{name}_mask"].append(np.asarray(view_b1.observed[name], dtype=bool))
|
|
native[f"{name}_coverage"].append(view_b1.coverage[name].mean(axis=1).astype(np.float16))
|
|
values, masks, cov = pooled[name]
|
|
rows[name].append(values.astype(np.float16))
|
|
rows[f"{name}_mask"].append(masks)
|
|
rows[f"{name}_coverage"].append(cov.mean(axis=1).astype(np.float32))
|
|
rows["text_posterior"].append(posterior[0].astype(np.float16))
|
|
rows["text_posterior_mask"].append(posterior[1])
|
|
rows["text_posterior_activity"].append(posterior[2].mean(axis=1).astype(np.float32))
|
|
bounds = np.column_stack((target_edges[:-1], target_edges[1:])).astype(np.float32)
|
|
centers = bounds.mean(axis=1)
|
|
rows["time_bounds_s"].append(bounds)
|
|
rows["time_s"].append(centers)
|
|
rows["progress"].append((centers / max(sample.duration_s, 1e-8)).astype(np.float32))
|
|
rows["duration_s"].append(float(sample.duration_s))
|
|
rows["valid_lengths"].append(EXPORT_BINS)
|
|
values = {
|
|
"id": sample_id, "sample_id": sample_id, "video_id": record["video_id"], "clip_id": record["clip_id"],
|
|
"raw_text": record["text"], "classification_labels": int(record["polarity"]),
|
|
"regression_labels": float(record["sentiment"]), "annotations": record["annotation"],
|
|
"label_consistent": bool(record["label_consistent"]), "source_video_sha256": source_hash,
|
|
}
|
|
for key, value in values.items():
|
|
rows[key].append(value)
|
|
|
|
args.output_dir.mkdir(parents=True, exist_ok=True)
|
|
arrays: dict[str, Any] = {}
|
|
string_keys = {"id", "sample_id", "video_id", "clip_id", "raw_text", "annotations", "source_video_sha256"}
|
|
for key, values in rows.items():
|
|
if key in string_keys:
|
|
arrays[key] = np.asarray(values, dtype=str)
|
|
elif key in {"duration_s", "regression_labels"}:
|
|
arrays[key] = np.asarray(values, dtype=np.float32)
|
|
elif key in {"classification_labels", "valid_lengths"}:
|
|
arrays[key] = np.asarray(values, dtype=np.int64)
|
|
elif key == "label_consistent":
|
|
arrays[key] = np.asarray(values, dtype=bool)
|
|
else:
|
|
arrays[key] = np.stack(values, axis=0)
|
|
|
|
dimensions = {"text": 768, "audio": len(cm.AUDIO_NAMES), "vision": len(cm.VISION_NAMES)}
|
|
names = {
|
|
"text": [f"bert_dim_{index:03d}" for index in range(768)],
|
|
"audio": list(cm.AUDIO_NAMES), "vision": list(cm.VISION_NAMES),
|
|
}
|
|
metadata = {
|
|
"schema": "q1-aligned50-v1", "sample_count": len(records), "export_bins": EXPORT_BINS,
|
|
"alignment": "Each clip is divided into 50 equal-width physical-time bins; values are duration/coverage-weighted means of Q1 B1 common-grid features.",
|
|
"native_grid_step_s": cm.GRID_STEP_S,
|
|
"time_semantics": "Seconds relative to clip start; time_bounds_s contains [left, right) for each bin.",
|
|
"text": "768-D BERT base uncased last-four-layer mean, word-level features projected by hard CTC Viterbi intervals with within-sample relative quality weighting.",
|
|
"text_posterior": "Alternative text projection uses fixed-transcript CTC forward-backward word occupancy. text_posterior_activity is expected occupancy mass, not physical coverage.",
|
|
"audio": "74-D: 40 log-Mel, 13 MFCC, 13 delta MFCC, 8 prosodic/spectral features.",
|
|
"vision": "35-D: 17 MediaPipe blendshape proxies, 6 head pose, 6 approximate gaze, 6 facial geometry; B1 uses SO(3) pose and normalized gaze pooling.",
|
|
"missing_values": "Values are zero where the matching boolean mask is false. Features are unnormalized; fit normalization on training groups only.",
|
|
"labels": {"classification_labels": {"0": "negative", "1": "neutral", "2": "positive"}, "regression_labels": "original continuous sentiment label"},
|
|
"split": "No train/validation/test split is supplied. Split by video_id to prevent source-video leakage.",
|
|
"dimensions": dimensions, "feature_names_file": "feature_names.json",
|
|
"source_run_manifest": str(args.run_manifest.relative_to(ROOT)) if args.run_manifest.is_relative_to(ROOT) else args.run_manifest.name,
|
|
}
|
|
pkl_path = args.output_dir / "aligned_50.pkl"
|
|
with pkl_path.open("wb") as stream:
|
|
pickle.dump({"metadata": metadata, "all": arrays}, stream, protocol=pickle.HIGHEST_PROTOCOL)
|
|
|
|
native_arrays: dict[str, Any] = {
|
|
"offsets": np.asarray(offsets, dtype=np.int32), "sample_indices": np.concatenate(native_ids).astype(np.int16),
|
|
"sample_id": np.asarray([record["sample_id"] for record in records], dtype=str),
|
|
"time_bounds_s": np.concatenate(native["time_bounds_s"], axis=0),
|
|
}
|
|
for key in ("text", "text_mask", "audio", "audio_mask", "vision", "vision_mask", "text_coverage", "audio_coverage", "vision_coverage"):
|
|
native_arrays[key] = np.concatenate(native[key], axis=0)
|
|
native_path = args.output_dir / "native_grid.npz"
|
|
np.savez_compressed(native_path, **native_arrays)
|
|
names_path = args.output_dir / "feature_names.json"
|
|
names_path.write_text(json.dumps(names, ensure_ascii=False, indent=2), encoding="utf-8")
|
|
|
|
result_dir = args.run_manifest.parent
|
|
for name in ("modality_summary.csv", "sample_alignment_summary.csv", "word_alignment_posterior.csv"):
|
|
source = result_dir / name
|
|
if source.exists():
|
|
shutil.copy2(source, args.output_dir / name)
|
|
summary_src = result_dir / "modality_summary.csv"
|
|
summary_path = args.output_dir / "sample_feature_manifest.csv"
|
|
modality_rows = []
|
|
with summary_src.open(encoding="utf-8-sig", newline="") as stream:
|
|
for row in csv.DictReader(stream):
|
|
row.update({
|
|
"source_extract_granularity": {
|
|
"text": "word vectors; hard CTC word intervals (posterior alternative uses 20 ms CTC frames)",
|
|
"audio": "10 ms feature hop", "vision": "5 Hz sampled frames (about 200 ms)",
|
|
}.get(row["modality"], ""),
|
|
"common_alignment_grid_s": str(cm.GRID_STEP_S), "training_export_bins": str(EXPORT_BINS),
|
|
"traceability_rule": "sample_id + source_video_sha256 + source time bounds",
|
|
})
|
|
modality_rows.append(row)
|
|
with summary_path.open("w", encoding="utf-8-sig", newline="") as stream:
|
|
writer = csv.DictWriter(stream, fieldnames=list(modality_rows[0]))
|
|
writer.writeheader()
|
|
writer.writerows(modality_rows)
|
|
|
|
files = {}
|
|
for path in (pkl_path, native_path, names_path, summary_path):
|
|
files[path.name] = {"bytes": path.stat().st_size, "sha256": sha256(path)}
|
|
feature_manifest = {
|
|
**metadata, "created_at_utc": datetime.now(timezone.utc).isoformat(),
|
|
"python": run_manifest.get("python"), "packages": run_manifest.get("packages", {}),
|
|
"models": run_manifest.get("models", {}), "label_file_sha256": label_hash,
|
|
"feature_names": names, "native_grid_step_count_total": int(offsets[-1]),
|
|
"native_grid_sample_offsets": offsets, "files": files,
|
|
"tables": {
|
|
"sample_feature_manifest.csv": "300 rows, one per sample and modality; records source/observed duration, dimension, granularity, grid length and source hash.",
|
|
"sample_alignment_summary.csv": "100 sample-level alignment and coverage summaries.",
|
|
"word_alignment_posterior.csv": "Word-level hard intervals and CTC posterior interval summaries.",
|
|
},
|
|
}
|
|
(args.output_dir / "feature_manifest.json").write_text(json.dumps(feature_manifest, ensure_ascii=False, indent=2), encoding="utf-8")
|
|
print(f"Exported {len(records)} samples: text={arrays['text'].shape}, audio={arrays['audio'].shape}, vision={arrays['vision'].shape}")
|
|
print(f"Native grid: {offsets[-1]} steps across {sum(durations):.3f} seconds")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|