Spaces:
Sleeping
Sleeping
File size: 4,921 Bytes
b0add2b | 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 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 | """Synthetic multivariate stream generator with ground truth labels."""
import asyncio
import itertools
import random
from collections.abc import AsyncIterator, Iterator
from typing import Any
from src.stream.models import Observation, Label
# Baseline means for Phase A
BASELINE_TEMP = 25.0
BASELINE_PRESSURE = 101.3
BASELINE_VIBRATION = 0.5
# Standard deviation for correlated noise
NOISE_SIGMA = 0.5
# Correlation: shared latent factor weight
LATENT_WEIGHT = 0.7
INDEPENDENT_WEIGHT = (1 - LATENT_WEIGHT**2) ** 0.5
def _correlated_noise(rng: random.Random) -> tuple[float, float, float]:
"""Generate correlated noise: one latent factor, then map to three sensors."""
z = rng.gauss(0, NOISE_SIGMA)
eps_t = rng.gauss(0, NOISE_SIGMA * INDEPENDENT_WEIGHT)
eps_p = rng.gauss(0, NOISE_SIGMA * INDEPENDENT_WEIGHT)
eps_v = rng.gauss(0, NOISE_SIGMA * INDEPENDENT_WEIGHT)
return (
LATENT_WEIGHT * z + eps_t,
LATENT_WEIGHT * z + eps_p,
LATENT_WEIGHT * z + eps_v,
)
def _generate_sync(
*,
phase_a_length: int = 300,
phase_b_length: int = 200,
phase_c_length: int = 300,
drift_magnitude: float = 0.3,
anomaly_rate: float = 0.05,
point_ratio: float = 0.6,
seed: int = 42,
) -> Iterator[Observation]:
rng = random.Random(seed)
t = 0.0
# Phase A: normal baseline
for i in range(phase_a_length):
nt, np_, nv = _correlated_noise(rng)
yield Observation(
timestamp=t,
temperature=BASELINE_TEMP + nt,
pressure=BASELINE_PRESSURE + np_,
vibration=BASELINE_VIBRATION + nv,
label="normal",
phase="A",
)
t += 1.0
# Phase B: linear drift
for i in range(phase_b_length):
alpha = (i + 1) / phase_b_length
mt = BASELINE_TEMP + drift_magnitude * alpha
mp = BASELINE_PRESSURE + drift_magnitude * alpha
mv = BASELINE_VIBRATION + drift_magnitude * alpha
nt, np_, nv = _correlated_noise(rng)
yield Observation(
timestamp=t,
temperature=mt + nt,
pressure=mp + np_,
vibration=mv + nv,
label="normal",
phase="B",
)
t += 1.0
# Post-drift means
end_temp = BASELINE_TEMP + drift_magnitude
end_pressure = BASELINE_PRESSURE + drift_magnitude
end_vibration = BASELINE_VIBRATION + drift_magnitude
# Phase C: post-drift + anomalies
# Anomaly injection: every ~1/anomaly_rate observations on average
obs_count = 0
anomaly_choice = itertools.cycle(["point_anomaly", "contextual_anomaly"])
next_anomaly_step = max(1, int(rng.expovariate(anomaly_rate)))
for i in range(phase_c_length):
obs_count += 1
nt, np_, nv = _correlated_noise(rng)
if obs_count >= next_anomaly_step:
# Inject anomaly; pick type by point_ratio
if rng.random() < point_ratio:
label: Label = "point_anomaly"
# Spike in one or more sensors (3-5 sigma)
spike = rng.uniform(3.0, 5.0) * NOISE_SIGMA
which = rng.sample([0, 1, 2], k=rng.randint(1, 3))
if 0 in which:
nt += spike * rng.choice([-1, 1])
if 1 in which:
np_ += spike * rng.choice([-1, 1])
if 2 in which:
nv += spike * rng.choice([-1, 1])
else:
label = "contextual_anomaly"
# High temp + low pressure, opposite to the normal positive correlation.
# Max per-feature z-score ~2.4, below the 2.5 Z-Score cutoff.
nt = 1.8 * NOISE_SIGMA
np_ = -1.8 * NOISE_SIGMA
nv = 0.0
next_anomaly_step = obs_count + max(1, int(rng.expovariate(anomaly_rate)))
else:
label = "normal"
yield Observation(
timestamp=t,
temperature=end_temp + nt,
pressure=end_pressure + np_,
vibration=end_vibration + nv,
label=label,
phase="C",
)
t += 1.0
async def generate_stream(
*,
phase_a_length: int = 300,
phase_b_length: int = 200,
phase_c_length: int = 300,
drift_magnitude: float = 0.3,
anomaly_rate: float = 0.05,
point_ratio: float = 0.6,
delay: float = 0.0,
seed: int = 42,
) -> AsyncIterator[Observation]:
"""Async generator yielding one observation at a time with optional delay."""
sync_gen = _generate_sync(
phase_a_length=phase_a_length,
phase_b_length=phase_b_length,
phase_c_length=phase_c_length,
drift_magnitude=drift_magnitude,
anomaly_rate=anomaly_rate,
point_ratio=point_ratio,
seed=seed,
)
for obs in sync_gen:
if delay > 0:
await asyncio.sleep(delay)
yield obs
|