File size: 4,961 Bytes
1616901 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 | #!/usr/bin/env python3
"""Run autoregressive FengWu-W2S inference and save forecast fields."""
from __future__ import annotations
import argparse
import json
import sys
from datetime import datetime, timedelta
from pathlib import Path
import numpy as np
import torch
import yaml
PROJECT_ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(PROJECT_ROOT))
from model.fengwu_w2s import FengWuW2S
from scripts.data_loader import ERA5WindowDataset, read_metadata, resolve_data_dir
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--config", default=str(PROJECT_ROOT / "conf/config.yaml"))
parser.add_argument("--data-dir")
parser.add_argument("--checkpoint")
parser.add_argument("--output-dir")
parser.add_argument("--split", choices=("train", "val", "test"), default="test")
parser.add_argument("--steps", type=int)
parser.add_argument("--limit", type=int, help="Limit windows for a quick smoke test")
parser.add_argument("--stochastic", action="store_true")
parser.add_argument("--seed", type=int, default=0)
return parser.parse_args()
def _resolve(path: str | Path, root: Path = PROJECT_ROOT) -> Path:
candidate = Path(path).expanduser()
return candidate if candidate.is_absolute() else (root / candidate).resolve()
def _build_model(config: dict) -> FengWuW2S:
model_cfg = config["model"]
data_cfg = config["data"]
return FengWuW2S(
in_channels=len(data_cfg["channels"]),
group_indices=data_cfg["groups"],
hidden_channels=int(model_cfg.get("hidden_channels", 32)),
latent_channels=int(model_cfg.get("latent_channels", 64)),
patch_size=int(model_cfg.get("patch_size", 4)),
num_blocks=int(model_cfg.get("num_blocks", 2)),
perturbation_scale=float(model_cfg.get("perturbation_scale", 0.01)),
)
def main() -> None:
args = parse_args()
with Path(args.config).open(encoding="utf-8") as source:
config = yaml.safe_load(source)
data_cfg = config["data"]
model_cfg = config["model"]
data_dir = resolve_data_dir(args.data_dir or data_cfg["data_dir"], PROJECT_ROOT)
checkpoint = _resolve(args.checkpoint or model_cfg.get("default_checkpoint", "./data/checkpoints/model_bak.pth"))
if not checkpoint.exists():
raise FileNotFoundError(f"Checkpoint not found: {checkpoint}; run scripts/train.py first")
split_years = {
"train": data_cfg["train_years"],
"val": data_cfg["val_years"],
"test": data_cfg["test_years"],
}[args.split]
steps = int(args.steps or model_cfg.get("inference_steps", 1))
dataset = ERA5WindowDataset(
data_dir=data_dir,
years=split_years,
channels=data_cfg["channels"],
input_steps=int(model_cfg.get("input_steps", 2)),
rollout_steps=max(1, steps),
)
metadata = read_metadata(data_dir, data_cfg["channels"])
means = torch.from_numpy(metadata["means"])
stds = torch.from_numpy(metadata["stds"])
device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
model = _build_model(config).to(device)
state = torch.load(checkpoint, map_location=device, weights_only=False)
model.load_state_dict(state["model_state_dict"])
model.eval()
output_dir = _resolve(args.output_dir or "./result/output")
output_dir.mkdir(parents=True, exist_ok=True)
index = []
limit = len(dataset) if args.limit is None else min(len(dataset), max(0, int(args.limit)))
for sample_index in range(limit):
inputs, _, timestamp = dataset[sample_index]
initial = inputs.unsqueeze(0).to(device)
with torch.no_grad():
prediction = model.rollout(initial, steps=steps, stochastic=args.stochastic, seed=args.seed + sample_index)
prediction = prediction.squeeze(0).cpu() * stds + means
year = timestamp[:4]
year_dir = output_dir / year
year_dir.mkdir(parents=True, exist_ok=True)
start = datetime.strptime(timestamp, "%Y%m%d%H")
for lead, field in enumerate(prediction):
valid_time = start + timedelta(hours=lead * int(metadata["time_step"]))
valid_timestamp = valid_time.strftime("%Y%m%d%H")
path = year_dir / f"{timestamp}_lead{lead:03d}.npy"
np.save(path, field.numpy().astype(np.float32))
index.append({
"source_timestamp": timestamp,
"valid_timestamp": valid_timestamp,
"lead": lead,
# Keep the manifest relocatable when the project is copied to
# a compute node or downloaded from ModelScope.
"path": str(path.relative_to(output_dir)),
})
(output_dir / "index.json").write_text(json.dumps(index, indent=2), encoding="utf-8")
print(f"Saved {len(index)} forecast fields under {output_dir}")
if __name__ == "__main__":
main()
|