Skip to content

Repository files navigation

SCALE: Streaming Causal-Aware Learning for Online Anomaly Detection

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.

Key Features

  • 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.

Requirements

pip install -r requirements.txt

Main 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

Repository Structure

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

Data Preparation

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.

scenarios_summary.csv Format

scenario_id,fault_category,num_windows,num_features,...
my_scenario_1,fault_A,2141,71,...

Configuration

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

Quick Start

1. Install dependencies

pip install -r requirements.txt

2. Prepare data

Place your datasets under data/processed_energy_storage/ following the structure described above. Update scenarios_summary.csv accordingly.

3. Run online detection

# 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 cpu

4. Visualize results

python visualize_results.py --results_dir results/my_scenario_1

How It Works

The framework operates in two sequential phases for each data stream:

Phase 1: Streaming Warm-up (warmup_phase_streaming)

  1. Normalize the first warmup_samples samples using fixed statistics (computed on the warm-up batch).
  2. Train Expert 0 — the shared Masked Causal VAE-TCN backbone plus the first style expert — using batch SGD for warmup_epochs epochs.
  3. Collect statistics — compute the SLR distribution and z_s statistics for Expert 0 to establish a drift detection baseline.
  4. Fix threshold — compute the anomaly score threshold from warm-up scores (fixed quantile mode).

Phase 2: Online Detection & Adaptation (online_detect_and_adapt_phase)

For each incoming sample at timestep t:

  1. Normalize using the current domain's fixed statistics.
  2. Route — compute SLR against all experts; select the expert with the minimum SLR.
  3. Score — compute reconstruction error (RE) and optionally prediction error (PRE); the final anomaly score is PRE if available, otherwise RE.
  4. Record — lock the prediction (prequential guarantee).
  5. 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.
  6. Create expert — when the drift buffer is full and samples are temporally continuous, freeze all existing experts and create a new one.

Output

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

Citation

If you use this code in your research, please cite our paper (details to be added upon publication).

License

This project is released for academic use. See LICENSE for details.

About

No description, website, or topics provided.

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages