This repository contains the core implementation of SCALE, an online Mixture-of-Experts (MoE) framework for time-series anomaly detection under concept drift. The framework supports fully streaming, prequential (test-then-train) evaluation without requiring a separate offline training set.
- Multi-Mask Anti-Contamination: A masked TCN encoder generates multiple random-mask views per sample, and an attention module aggregates them into robust latent representations.
- Adaptive Expert Pool: Experts are created dynamically upon domain drift detection, each memorizing one "normal domain" distribution.
- SLR-Based Drift Detection: Sequence Likelihood Regret (SLR) measures how well a sample fits each expert; persistent high SLR triggers new expert creation.
- Koopman Operator Predictor: A linear Koopman operator predicts the next causal feature from history, providing an auxiliary anomaly signal (PRE).
- Strict Prequential Evaluation: Predictions are locked before any model update, preventing information leakage.
pip install -r requirements.txtMain dependencies:
- Python >= 3.8
- PyTorch >= 1.9.0
- NumPy >= 1.19.0
- scikit-learn >= 0.23.0
- PyYAML >= 5.3.0
- pandas >= 1.1.0
- tqdm
scale_core_3/
├── README.md
├── requirements.txt
├── test_energy_storage.py # Main entry: streaming anomaly detection
├── visualize_results.py # Visualization of detection results
├── core/ # Core MoE detector logic
│ ├── online_moe_detector.py # Main detector class
│ ├── masked_style_expert.py # Per-domain expert module
│ ├── expert_similarity.py # Expert similarity & merge-cost evaluation
│ ├── _online_adaptation.py # Prequential online detect-and-adapt loop
│ ├── _online_methods.py # Method binding for streaming mode
│ ├── _training.py # Warm-up training
│ ├── _slr_computation.py # SLR computation and expert routing
│ └── _loss_computation.py # Loss functions (diversity, HSIC)
├── models/
│ ├── masked_causal_vae_tcn.py # Masked Causal VAE-TCN (shared backbone)
│ └── koopman_predictor.py # Koopman operator predictor
├── utils/
│ ├── data_utils.py # Data loading & normalization
│ ├── config_utils.py # YAML config loading
│ ├── online_normalizer.py # Online Welford normalization
│ ├── buffer_system.py # Drift candidate buffer
│ ├── stream_history_buffer.py # Sliding history window
│ ├── update_gate.py # Update gate (anti-contamination)
│ ├── expert_statistics.py # Per-expert incremental statistics
│ ├── fixed_quantile_threshold.py # Fixed anomaly threshold
│ ├── global_adaptive_threshold.py # Adaptive anomaly threshold
│ ├── sequential_data_loader.py # Sequential DataLoader for predictors
│ └── metrics.py # Evaluation metrics (F1, AUC, VUS-PR, etc.)
├── config/
│ └── exo_standard_single_2.yaml # Example hyperparameter configuration
└── data/
└── processed_energy_storage/ # Place your datasets here
└── scenarios_summary.csv # Dataset index file
The framework expects pre-processed .npy data files with the following structure:
data/processed_energy_storage/
├── scenarios_summary.csv
└── <fault_category>/
├── <scenario_id>_test_with_anomalies.npy # shape: (N, D)
├── <scenario_id>_test_anomaly_label.npy # shape: (N,), values: {0, 1}
└── <scenario_id>_test_domain_label.npy # shape: (N,), optional
test_with_anomalies.npy: flattened time-series windows(N, num_features × window_size).test_anomaly_label.npy: binary labels,0= normal,1= anomaly.scenarios_summary.csv: index file listing all available scenarios.
scenario_id,fault_category,num_windows,num_features,...
my_scenario_1,fault_A,2141,71,...Hyperparameters are managed via YAML files located in config/. See config/exo_standard_single_2.yaml for a fully annotated example.
Key configuration sections:
| Section | Key Parameters | Description |
|---|---|---|
model |
latent_dim_c, latent_dim_s |
Causal and style latent dimensions |
masked_tcn |
num_masks, mask_dropout |
Number of random masks per forward pass |
warmup |
warmup_samples, warmup_epochs |
Streaming warm-up configuration |
online_training |
buffer_size, slr_quantile, slr_margin |
Drift detection parameters |
predictor |
predictor_type, history_length |
Koopman predictor configuration |
system |
device, seed |
Hardware and reproducibility |
pip install -r requirements.txtPlace your datasets under data/processed_energy_storage/ following the structure described above. Update scenarios_summary.csv accordingly.
# Run on all datasets listed in scenarios_summary.csv
python test_energy_storage.py
# Run on a specific dataset (skip visualization for speed)
python test_energy_storage.py --dataset my_scenario_1 --skip_visualization
# Adjust warm-up sample count
python test_energy_storage.py --warmup_samples 400 --skip_visualization
# Use CPU instead of GPU
python test_energy_storage.py --device cpupython visualize_results.py --results_dir results/my_scenario_1The framework operates in two sequential phases for each data stream:
- Normalize the first
warmup_samplessamples using fixed statistics (computed on the warm-up batch). - Train Expert 0 — the shared Masked Causal VAE-TCN backbone plus the first style expert — using batch SGD for
warmup_epochsepochs. - Collect statistics — compute the SLR distribution and z_s statistics for Expert 0 to establish a drift detection baseline.
- Fix threshold — compute the anomaly score threshold from warm-up scores (fixed quantile mode).
For each incoming sample at timestep t:
- Normalize using the current domain's fixed statistics.
- Route — compute SLR against all experts; select the expert with the minimum SLR.
- Score — compute reconstruction error (RE) and optionally prediction error (PRE); the final anomaly score is PRE if available, otherwise RE.
- Record — lock the prediction (prequential guarantee).
- Decide — categorize the sample:
normal: SLR is low and score is below threshold → update the current expert.drift: SLR is high → add to the drift buffer.anomaly: score exceeds threshold → do not update.
- Create expert — when the drift buffer is full and samples are temporally continuous, freeze all existing experts and create a new one.
Results are saved under results/<scenario_id>/:
results/my_scenario_1/
├── online_results.npz # Per-timestep detection results
├── metrics.json # Summary metrics (AUC-ROC, F1-T, etc.)
└── figures/ # Visualization plots
Key fields in online_results.npz:
| Field | Description |
|---|---|
final_score |
Anomaly score at each timestep |
expert_id |
Expert assigned at each timestep |
sample_type |
'normal', 'anomaly', or 'drift' |
is_drift |
Whether domain drift was detected |
new_expert_created |
Whether a new expert was created at this timestep |
threshold |
Decision threshold at each timestep |
If you use this code in your research, please cite our paper (details to be added upon publication).
This project is released for academic use. See LICENSE for details.