diff --git a/docs/api/processors.rst b/docs/api/processors.rst index a06e3c955..632433622 100644 --- a/docs/api/processors.rst +++ b/docs/api/processors.rst @@ -299,6 +299,64 @@ Common string keys for automatic processor selection: - ``"raw"``: For raw/unprocessed data - ``"graph"``: For knowledge graph subgraphs +Training-only time-series normalization +-------------------------------------- + +``TimeseriesProcessor`` and ``TemporalTimeseriesProcessor`` accept +``normalize_strategy="standard"`` for per-feature z-score normalization. +The default, ``None``, keeps existing behavior. Min–max scaling is not supported. + +``fit()`` computes the training mean and population standard deviation +(``ddof=0``) **after resampling and imputation**. Every resampled time step has +equal weight, including filled gaps; longer sequences therefore contribute +more observations. Constant features use a scale of 1. ``process()`` applies +these fixed statistics without updating them. The temporal processor normalizes +only ``"value"``; its ``"time"`` tensor remains in elapsed hours. + +.. important:: + + Split raw samples by patient before fitting normalization. Calling + ``set_task()`` with normalization enabled on the full dataset and then + splitting the resulting samples leaks validation/test statistics into training. + +Create the training dataset first, then explicitly transfer its fitted processors: + +.. code-block:: python + + from pyhealth.datasets import create_sample_dataset + + schema = {"vitals": ("timeseries", {"normalize_strategy": "standard"})} + train = create_sample_dataset( + samples=train_samples, input_schema=schema, output_schema={} + ) + validation = create_sample_dataset( + samples=validation_samples, + input_schema=schema, + output_schema={}, + input_processors=train.input_processors, + ) + +Use the same transfer for the test split. For temporal data, replace +``"timeseries"`` with ``"temporal_timeseries"``. Passing a pre-fitted processor +instance in ``input_schema`` alone does **not** prevent refitting; use the +``input_processors`` argument instead. The processor cannot infer which samples +are training data, so the caller must supply the intended training partition. + +Normalization requires a successful ``fit()``; otherwise processing raises +``RuntimeError``. Refitting replaces all previous statistics, and failed fitting +leaves normalization unfitted. Missing or ``None`` training fields are skipped; +no usable training samples, malformed present samples, and infinite values raise +``ValueError``. NaNs are handled by the existing imputation step, including its +zero initialization for leading or entirely missing features. Resampled and +imputed values must be finite. Existing input shapes and imputation behavior +are preserved: the ordinary processor expects ``(T, F)`` values, while the +temporal processor also accepts ``(T,)``. + +Fitted statistics are saved with the existing processor metadata and reused by +cached datasets; no separate statistics file is necessary. See +``examples/timeseries_normalization.py`` for a complete synthetic example that +can be run with ``pixi run -e test python examples/timeseries_normalization.py``. + Writing Custom FeatureProcessors --------------------------------- @@ -496,4 +554,4 @@ API Reference processors/pyhealth.processors.MultiHotProcessor processors/pyhealth.processors.StageNetProcessor processors/pyhealth.processors.StageNetTensorProcessor - processors/pyhealth.processors.GraphProcessor \ No newline at end of file + processors/pyhealth.processors.GraphProcessor diff --git a/examples/timeseries_normalization.py b/examples/timeseries_normalization.py new file mode 100644 index 000000000..ceff8600b --- /dev/null +++ b/examples/timeseries_normalization.py @@ -0,0 +1,59 @@ +"""Fit z-score normalization on training patients and reuse it on held-out data. + +Run with: pixi run -e test python examples/timeseries_normalization.py +All samples are synthetic; no downloads or model training are required. +""" + +from datetime import datetime, timedelta + +import numpy as np + +from pyhealth.datasets import create_sample_dataset + + +def main() -> None: + """Demonstrate ordinary and temporal normalization without split leakage.""" + start = datetime(2026, 1, 1) + raw_samples = [ + { + "patient_id": "train-patient", + "vitals": ([start, start + timedelta(hours=1)], np.array([[10.0], [30.0]])), + }, + {"patient_id": "validation-patient", "vitals": ([start], np.array([[40.0]]))}, + {"patient_id": "test-patient", "vitals": ([start], np.array([[1000.0]]))}, + ] + # Split raw samples by patient BEFORE fitting any normalization statistics. + split_patients = { + "train": {"train-patient"}, + "validation": {"validation-patient"}, + "test": {"test-patient"}, + } + splits = { + name: [s for s in raw_samples if s["patient_id"] in patients] + for name, patients in split_patients.items() + } + + for alias in ("timeseries", "temporal_timeseries"): + schema = {"vitals": (alias, {"normalize_strategy": "standard"})} + train = create_sample_dataset( + samples=splits["train"], input_schema=schema, output_schema={} + ) + for name, expected in (("validation", 2.0), ("test", 98.0)): + held_out = create_sample_dataset( + samples=splits[name], + input_schema=schema, + output_schema={}, + # Supplying processors here prevents fitting them on held-out data. + input_processors=train.input_processors, + ) + value = held_out[0]["vitals"] + if isinstance(value, dict): + np.testing.assert_array_equal(value["time"], [0.0]) + value = value["value"] + # Training mean=20, scale=10. The test outlier does not change them. + np.testing.assert_array_equal(value, [[expected]]) + print(f"{alias}: {name} normalized values = {value.tolist()}") + + +if __name__ == "__main__": + main() diff --git a/pyhealth/processors/temporal_timeseries_processor.py b/pyhealth/processors/temporal_timeseries_processor.py index b38a84295..8de6eb43e 100644 --- a/pyhealth/processors/temporal_timeseries_processor.py +++ b/pyhealth/processors/temporal_timeseries_processor.py @@ -8,14 +8,20 @@ ``UnifiedMultimodalEmbeddingModel`` can sort and align events across modalities on a shared timeline. """ + from datetime import datetime, timedelta -from typing import Any, Iterable, Dict, List, Tuple +from typing import Any, Iterable, Dict, List, Literal, Tuple import numpy as np import torch from . import register_processor from .base_processor import ModalityType, TemporalFeatureProcessor +from .timeseries_processor import ( + _fit_standardization, + _standardize, + _validate_normalization_input, +) @register_processor("temporal_timeseries") @@ -38,6 +44,10 @@ class TemporalTimeseriesProcessor(TemporalFeatureProcessor): Args: sampling_rate: Uniform re-sampling interval. Defaults to 1 hour. impute_strategy: Currently only ``"forward_fill"`` is supported. + normalize_strategy: ``None`` (default) or ``"standard"``. Z-scores use + training-only statistics after resampling and imputation, weighting + each time step equally. Constant features use scale 1. Only the + ``"value"`` tensor is normalized; timestamps are unchanged. Example:: @@ -48,21 +58,57 @@ class TemporalTimeseriesProcessor(TemporalFeatureProcessor): out = proc.process_temporal((ts, val)) # out["value"].shape → (5, 2) ← 5 two-hour steps over 8 h # out["time"].shape → (5,) ← [0., 2., 4., 6., 8.] hours + + Examples: + >>> times = [datetime(2026, 1, 1), datetime(2026, 1, 1, 1)] + >>> values = np.array([[10.0], [30.0]]) + >>> processor = TemporalTimeseriesProcessor(normalize_strategy="standard") + >>> processor.fit([{"vitals": (times, values)}], "vitals") + >>> output = processor.process((times, values)) + >>> output["value"].tolist() + [[-1.0], [1.0]] + >>> output["time"].tolist() + [0.0, 1.0] """ def __init__( self, sampling_rate: timedelta = timedelta(hours=1), impute_strategy: str = "forward_fill", + normalize_strategy: Literal["standard"] | None = None, ): + if normalize_strategy not in (None, "standard"): + raise ValueError("normalize_strategy must be None or 'standard'.") + if normalize_strategy is not None and sampling_rate <= timedelta(0): + raise ValueError("sampling_rate must be positive for normalization.") self.sampling_rate = sampling_rate self.impute_strategy = impute_strategy self.n_features: int | None = None + self.normalize_strategy = normalize_strategy + self._normalization_mean: list[float] | None = None + self._normalization_scale: list[float] | None = None # ── FeatureProcessor interface ───────────────────────────────────────── def fit(self, samples: Iterable[Dict[str, Any]], field: str) -> None: - """Infer feature dimension from the first valid sample.""" + """Infer feature count and optionally fit training-only z-score statistics. + + Statistics use population standard deviation (ddof=0) after filling. + Each fit replaces previous statistics, and a failed fit leaves + normalization unfitted. Missing/None fields are skipped; normalization + requires at least one usable training sample. + """ + self._normalization_mean = None + self._normalization_scale = None + if getattr(self, "normalize_strategy", None) == "standard": + self.n_features = None + ( + self.n_features, + self._normalization_mean, + self._normalization_scale, + ) = _fit_standardization(samples, field, self._resample_and_impute) + return + for sample in samples: if field in sample and sample[field] is not None: _, values = sample[field] @@ -83,7 +129,30 @@ def process(self, value: Tuple[List[datetime], np.ndarray]) -> dict: Returns: ``{"value": FloatTensor (S, F), "time": FloatTensor (S,)}`` + + Raises: + RuntimeError: If normalization is enabled but fitting has not succeeded. + ValueError: If the series is invalid for normalization. """ + sampled_values = self._resample_and_impute(value) + if getattr(self, "normalize_strategy", None) == "standard": + sampled_values = _standardize( + sampled_values, self._normalization_mean, self._normalization_scale + ) + + hours_per_step = self.sampling_rate.total_seconds() / 3600.0 + time_hours = np.array( + [i * hours_per_step for i in range(len(sampled_values))], dtype=np.float32 + ) + return { + "value": torch.tensor(sampled_values, dtype=torch.float32), + "time": torch.tensor(time_hours, dtype=torch.float32), + } + + def _resample_and_impute( + self, value: Tuple[List[datetime], np.ndarray] + ) -> np.ndarray: + """Produce the same unnormalized grid for fitting and processing.""" timestamps, values = value if len(timestamps) == 0: @@ -93,6 +162,11 @@ def process(self, value: Tuple[List[datetime], np.ndarray]) -> dict: if values.ndim == 1: values = values[:, None] # (T,) → (T, 1) + if getattr(self, "normalize_strategy", None) == "standard": + _validate_normalization_input( + timestamps, values, self.sampling_rate, self.n_features + ) + num_features = values.shape[1] start_time = timestamps[0] end_time = timestamps[-1] @@ -114,16 +188,7 @@ def process(self, value: Tuple[List[datetime], np.ndarray]) -> dict: else: sampled_values[i, f] = last - # Build time tensor (hours from first observation) - hours_per_step = self.sampling_rate.total_seconds() / 3600.0 - time_hours = np.array( - [i * hours_per_step for i in range(total_steps)], dtype=np.float32 - ) - - return { - "value": torch.tensor(sampled_values, dtype=torch.float32), - "time": torch.tensor(time_hours, dtype=torch.float32), - } + return sampled_values # process_temporal delegates to process (already returns dict) def process_temporal(self, value) -> dict: @@ -157,8 +222,10 @@ def size(self) -> int | None: return self.n_features def __repr__(self) -> str: + strategy = getattr(self, "normalize_strategy", None) + normalization = f", normalize_strategy={strategy!r}" if strategy else "" return ( f"TemporalTimeseriesProcessor(" f"sampling_rate={self.sampling_rate}, " - f"n_features={self.n_features})" + f"n_features={self.n_features}{normalization})" ) diff --git a/pyhealth/processors/timeseries_processor.py b/pyhealth/processors/timeseries_processor.py index e761e8035..f0369ce5f 100644 --- a/pyhealth/processors/timeseries_processor.py +++ b/pyhealth/processors/timeseries_processor.py @@ -1,13 +1,77 @@ from datetime import datetime, timedelta -from typing import Any, List, Tuple +from collections.abc import Callable, Iterable +from typing import Any, List, Literal, Tuple import numpy as np import torch +from sklearn.preprocessing import StandardScaler from . import register_processor from .base_processor import FeatureProcessor +def _validate_normalization_input( + timestamps: List[datetime], + values: np.ndarray, + sampling_rate: timedelta, + n_features: int | None, +) -> None: + """Validate a time series before fitting or applying normalization.""" + if sampling_rate <= timedelta(0): + raise ValueError("sampling_rate must be positive for normalization.") + if len(timestamps) == 0: + raise ValueError("Timestamps list is empty.") + if values.ndim != 2 or values.shape[1] == 0: + raise ValueError("Normalization requires values with shape (T, F), F > 0.") + if len(timestamps) != len(values): + raise ValueError("Timestamps and values must have the same length.") + if any(left > right for left, right in zip(timestamps, timestamps[1:])): + raise ValueError("Timestamps must be in nondecreasing order.") + if n_features is not None and values.shape[1] != n_features: + raise ValueError("Feature count does not match the fitted training data.") + if np.isinf(values).any(): + raise ValueError("Infinite values are not supported for normalization.") + + +def _fit_standardization( + samples: Iterable[dict[str, Any]], + field: str, + preprocess: Callable[[Any], np.ndarray], +) -> tuple[int, list[float], list[float]]: + """Accumulate population statistics one processed training sample at a time.""" + scaler = StandardScaler() + fitted = False + for sample in samples: + if field not in sample or sample[field] is None: + continue + values = preprocess(sample[field]) + if not np.isfinite(values).all(): + raise ValueError("Resampled and imputed training values must be finite.") + scaler.partial_fit(values) + fitted = True + if not fitted: + raise ValueError(f"No training samples contain usable data for {field!r}.") + if not np.isfinite(scaler.mean_).all() or not np.isfinite(scaler.scale_).all(): + raise ValueError( + "Training values produced non-finite normalization statistics." + ) + # Persist explicit numbers: BaseDataset fingerprints vars(processor) as JSON. + return scaler.n_features_in_, scaler.mean_.tolist(), scaler.scale_.tolist() + + +def _standardize( + values: np.ndarray, + mean: list[float] | None, + scale: list[float] | None, +) -> np.ndarray: + """Apply fixed training statistics without changing processor state.""" + if mean is None or scale is None: + raise RuntimeError("Call fit() on training samples before normalization.") + if not np.isfinite(values).all(): + raise ValueError("Resampled and imputed values must be finite.") + return (values - np.asarray(mean)) / np.asarray(scale) + + @register_processor("timeseries") class TimeseriesProcessor(FeatureProcessor): """ @@ -20,28 +84,70 @@ class TimeseriesProcessor(FeatureProcessor): Processing: 1. Uniform sampling at fixed intervals. 2. Imputation for missing values. + 3. Optional z-score normalization using training-set statistics. Output: - torch.Tensor of shape (S, F), where S is the number of sampled time steps. + + Args: + sampling_rate: Uniform sampling interval; defaults to one hour. + impute_strategy: ``"forward_fill"`` (default) or ``"zero"``. + normalize_strategy: ``None`` (default) or ``"standard"`` for z-scores. + With normalization enabled, call ``fit()`` on training samples only. + Statistics include all resampled, imputed time steps with equal + weight. Constant features use scale 1. Validation/test samples must + reuse the fitted processor rather than refitting it. + + Examples: + >>> times = [datetime(2026, 1, 1), datetime(2026, 1, 1, 1)] + >>> values = np.array([[10.0], [30.0]]) + >>> processor = TimeseriesProcessor(normalize_strategy="standard") + >>> processor.fit([{"vitals": (times, values)}], "vitals") + >>> processor.process((times, values)).tolist() + [[-1.0], [1.0]] """ def __init__( self, sampling_rate: timedelta = timedelta(hours=1), impute_strategy: str = "forward_fill", + normalize_strategy: Literal["standard"] | None = None, ): + if normalize_strategy not in (None, "standard"): + raise ValueError("normalize_strategy must be None or 'standard'.") + if normalize_strategy is not None and sampling_rate <= timedelta(0): + raise ValueError("sampling_rate must be positive for normalization.") # Configurable sampling rate and imputation method self.sampling_rate = sampling_rate self.impute_strategy = impute_strategy self.n_features = None + self.normalize_strategy = normalize_strategy + self._normalization_mean: list[float] | None = None + self._normalization_scale: list[float] | None = None def fit(self, samples: Any, field: str) -> None: - """Fit the processor by determining n_features from the first valid sample. + """Infer feature count and optionally fit training-only z-score statistics. + + Normalization uses the population standard deviation (ddof=0) across + resampled, imputed time steps. Each fit replaces previous statistics; + a failed fit leaves normalization unfitted. Missing/None fields are + skipped, but an empty training collection cannot fit normalization. Args: samples: Iterable of sample dictionaries. field: The field name to extract from samples. """ + self._normalization_mean = None + self._normalization_scale = None + if getattr(self, "normalize_strategy", None) == "standard": + self.n_features = None + ( + self.n_features, + self._normalization_mean, + self._normalization_scale, + ) = _fit_standardization(samples, field, self._resample_and_impute) + return + # Extract n_features from the first valid sample without full processing for sample in samples: if field in sample and sample[field] is not None: @@ -55,12 +161,37 @@ def fit(self, samples: Any, field: str) -> None: break def process(self, value: Tuple[List[datetime], np.ndarray]) -> torch.Tensor: + """Resample, impute, and optionally apply fixed training statistics. + + Raises: + RuntimeError: If normalization is enabled but fitting has not succeeded. + ValueError: If the series is invalid for normalization. + """ + sampled_values = self._resample_and_impute(value) + if getattr(self, "normalize_strategy", None) == "standard": + sampled_values = _standardize( + sampled_values, self._normalization_mean, self._normalization_scale + ) + elif self.n_features is None: + self.n_features = sampled_values.shape[1] + return torch.tensor(sampled_values, dtype=torch.float) + + def _resample_and_impute( + self, value: Tuple[List[datetime], np.ndarray] + ) -> np.ndarray: + """Produce the same unnormalized grid for fitting and processing.""" timestamps, values = value if len(timestamps) == 0: raise ValueError("Timestamps list is empty.") - values = np.asarray(values) + if getattr(self, "normalize_strategy", None) == "standard": + values = np.asarray(values, dtype=np.float64) + _validate_normalization_input( + timestamps, values, self.sampling_rate, self.n_features + ) + else: + values = np.asarray(values) num_features = values.shape[1] # Step 1: Uniform sampling @@ -68,9 +199,6 @@ def process(self, value: Tuple[List[datetime], np.ndarray]) -> torch.Tensor: end_time = timestamps[-1] total_steps = int((end_time - start_time) / self.sampling_rate) + 1 - sampled_times = [ - start_time + i * self.sampling_rate for i in range(total_steps) - ] sampled_values = np.full((total_steps, num_features), np.nan) # Map original timestamps to indices in the sampled grid @@ -93,10 +221,7 @@ def process(self, value: Tuple[List[datetime], np.ndarray]) -> torch.Tensor: else: raise ValueError(f"Unsupported imputation strategy: {self.impute_strategy}") - if self.n_features is None: - self.n_features = sampled_values.shape[1] - - return torch.tensor(sampled_values, dtype=torch.float) + return sampled_values def size(self): # Size equals number of features, unknown until first process @@ -118,7 +243,9 @@ def spatial(self) -> tuple[bool, ...]: return (True, False) def __repr__(self): + strategy = getattr(self, "normalize_strategy", None) + normalization = f", normalize_strategy={strategy!r}" if strategy else "" return ( f"TimeSeriesProcessor(sampling_rate={self.sampling_rate}, " - f"impute_strategy='{self.impute_strategy}')" + f"impute_strategy='{self.impute_strategy}'{normalization})" ) diff --git a/tests/core/test_timeseries_normalization.py b/tests/core/test_timeseries_normalization.py new file mode 100644 index 000000000..4f59da516 --- /dev/null +++ b/tests/core/test_timeseries_normalization.py @@ -0,0 +1,299 @@ +"""Training-only normalization contracts for both time-series processors.""" + +import copy +from datetime import datetime, timedelta +import json +from pathlib import Path +import pickle +import tempfile +import unittest + +import numpy as np +import torch + +from pyhealth.datasets import create_sample_dataset +from pyhealth.datasets.sample_dataset import SampleBuilder +from pyhealth.processors import TemporalTimeseriesProcessor, TimeseriesProcessor + + +PROCESSORS = (TimeseriesProcessor, TemporalTimeseriesProcessor) + + +def series(values, hours=None): + values = np.asarray(values, dtype=float) + if hours is None: + hours = range(len(values)) + timestamps = [datetime(2026, 1, 1) + timedelta(hours=h) for h in hours] + return timestamps, values + + +def tensor(processor, value): + result = processor.process(value) + return result["value"] if isinstance(result, dict) else result + + +class TestTimeseriesNormalization(unittest.TestCase): + def test_default_preserves_resampling_and_filling(self): + value = series([[np.nan, 10], [4, np.nan]], hours=[0, 2]) + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls() + result = tensor(proc, value) + self.assertEqual(result.dtype, torch.float32) + np.testing.assert_array_equal(result, [[0, 10], [0, 10], [4, 10]]) + + def test_standard_uses_all_resampled_steps_and_generator_once(self): + # Filled training values are [10, 10, 30, 50]; each step has equal weight. + samples = [ + {"signal": series([[10, 7], [30, 7]], hours=[0, 2])}, + {"signal": series([[50, 7]])}, + ] + expected = (np.array([10, 10, 30, 50]) - 25) / np.sqrt(275) + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + proc.fit(iter(samples), "signal") + result = torch.cat([tensor(proc, s["signal"]) for s in samples]) + np.testing.assert_allclose(result[:, 0], expected, rtol=1e-6) + np.testing.assert_array_equal(result[:, 1], 0) + self.assertEqual(proc.size(), 2) + self.assertEqual(result.dtype, torch.float32) + + def test_zero_imputation_precedes_normalization(self): + proc = TimeseriesProcessor( + impute_strategy="zero", normalize_strategy="standard" + ) + value = series([[10], [30]], hours=[0, 2]) + proc.fit([{"signal": value}], "signal") + expected = np.array([[10], [0], [30]], dtype=float) + expected = (expected - expected.mean()) / expected.std() + np.testing.assert_allclose(tensor(proc, value), expected, rtol=1e-6) + + def test_temporal_time_is_unchanged_and_1d_is_supported(self): + value = series([10, 30], hours=[0, 4]) + proc = TemporalTimeseriesProcessor( + sampling_rate=timedelta(hours=2), normalize_strategy="standard" + ) + proc.fit([{"signal": value}], "signal") + result = proc.process_temporal(value) + self.assertEqual(result["value"].shape, (3, 1)) + np.testing.assert_array_equal(result["time"], [0, 2, 4]) + self.assertEqual(result["time"].dtype, torch.float32) + np.testing.assert_allclose( + result["value"][:, 0], + [-1 / np.sqrt(2), -1 / np.sqrt(2), np.sqrt(2)], + rtol=1e-6, + ) + + def test_ordinary_processor_rejects_1d_values_cleanly(self): + proc = TimeseriesProcessor(normalize_strategy="standard") + with self.assertRaisesRegex(ValueError, "shape"): + proc.fit([{"signal": series([10, 30])}], "signal") + + def test_constant_single_step_and_all_missing_features(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + proc.fit([{"signal": series([[7, np.nan]])}], "signal") + np.testing.assert_array_equal(tensor(proc, series([[7, np.nan]])), 0) + np.testing.assert_array_equal(tensor(proc, series([[9, 3]])), [[2, 3]]) + + def test_processing_does_not_modify_inputs_or_fitted_state(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + proc.fit([{"signal": series([[10], [30]])}], "signal") + before = copy.deepcopy(vars(proc)) + value = series([[1000], [np.nan]]) + original = value[1].copy() + np.testing.assert_array_equal(tensor(proc, value), [[98], [98]]) + self.assertEqual(vars(proc), before) + np.testing.assert_array_equal(value[1], original) + + def test_invalid_strategy(self): + for cls in PROCESSORS: + for strategy in ("minmax", "unknown", True): + with self.subTest(processor=cls.__name__, strategy=strategy): + with self.assertRaisesRegex(ValueError, "normalize_strategy"): + cls(normalize_strategy=strategy) + + def test_fit_is_required_only_when_normalization_enabled(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + value = series([[10], [30]]) + np.testing.assert_array_equal(tensor(cls(), value), value[1]) + with self.assertRaisesRegex(RuntimeError, "fit"): + tensor(cls(normalize_strategy="standard"), value) + + def test_missing_fields_are_skipped_but_no_data_raises(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + proc.fit( + [{}, {"signal": None}, {"signal": series([[1], [3]])}], "signal" + ) + np.testing.assert_array_equal( + tensor(proc, series([[1], [3]])), [[-1], [1]] + ) + for samples in ([], [{}, {"signal": None}]): + with self.assertRaisesRegex(ValueError, "training"): + proc.fit(samples, "signal") + with self.assertRaises(RuntimeError): + tensor(proc, series([[1]])) + + def test_refit_replaces_statistics_and_failed_refit_clears_them(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + proc.fit([{"signal": series([[10], [30]])}], "signal") + proc.fit([{"signal": series([[100], [300]])}], "signal") + np.testing.assert_array_equal( + tensor(proc, series([[100], [300]])), [[-1], [1]] + ) + with self.assertRaises(ValueError): + proc.fit( + [ + {"signal": series([[1]])}, + {"signal": series([[np.inf]])}, + ], + "signal", + ) + with self.assertRaises(RuntimeError): + tensor(proc, series([[100]])) + + def test_invalid_present_samples_rejected_during_fit_and_process(self): + bad_values = [ + series(np.empty((0, 2))), + series([[1, 2]], hours=[0, 1]), + series([[1, 2], [3, 4]], hours=[1, 0]), + series([[[1, 2]]]), + series(np.empty((1, 0))), + series([[np.inf, 1]]), + series([[-np.inf, 1]]), + ] + for cls in PROCESSORS: + for value in bad_values: + with self.subTest(processor=cls.__name__, shape=value[1].shape): + proc = cls(normalize_strategy="standard") + with self.assertRaises(ValueError): + proc.fit([{"signal": value}], "signal") + proc.fit([{"signal": series([[1, 2], [3, 4]])}], "signal") + with self.assertRaises(ValueError): + proc.process(value) + + def test_feature_count_must_match(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + with self.assertRaises(ValueError): + proc.fit( + [{"signal": series([[1, 2]])}, {"signal": series([[1]])}], + "signal", + ) + proc.fit([{"signal": series([[1, 2]])}], "signal") + with self.assertRaises(ValueError): + tensor(proc, series([[1]])) + + def test_nonpositive_sampling_rate_rejected_when_enabled(self): + for cls in PROCESSORS: + for rate in (timedelta(0), timedelta(hours=-1)): + with self.subTest(processor=cls.__name__, rate=rate): + with self.assertRaisesRegex(ValueError, "sampling_rate"): + cls(sampling_rate=rate, normalize_strategy="standard") + + def test_duplicate_timestamps_keep_existing_last_value_rule(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls(normalize_strategy="standard") + value = series([[1], [10], [30]], hours=[0, 0, 1]) + proc.fit([{"signal": value}], "signal") + np.testing.assert_array_equal(tensor(proc, value), [[-1], [1]]) + + def test_legacy_pickle_without_normalization_fields(self): + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + proc = cls() + proc.__dict__ = { + "sampling_rate": timedelta(hours=1), + "impute_strategy": "forward_fill", + "n_features": 1, + } + restored = pickle.loads(pickle.dumps(proc)) + np.testing.assert_array_equal( + tensor(restored, series([[10], [30]])), [[10], [30]] + ) + restored.fit([{"signal": series([[20]])}], "signal") + np.testing.assert_array_equal(tensor(restored, series([[20]])), [[20]]) + + +class TestNormalizationIntegration(unittest.TestCase): + def test_schema_instance_fingerprint_includes_normalization(self): + # Task schemas containing instances use default=str in BaseDataset. + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + plain = json.dumps({"signal": cls()}, sort_keys=True, default=str) + normalized = json.dumps( + {"signal": cls(normalize_strategy="standard")}, + sort_keys=True, + default=str, + ) + self.assertNotEqual(plain, normalized) + old = cls() + del old.normalize_strategy + self.assertEqual(repr(old), repr(cls())) + + def test_training_processor_transfer_does_not_refit_on_held_out_data(self): + for alias in ("timeseries", "temporal_timeseries"): + with self.subTest(processor=alias): + schema = {"signal": (alias, {"normalize_strategy": "standard"})} + train = create_sample_dataset( + samples=[{"patient_id": "train", "signal": series([[10], [30]])}], + input_schema=schema, + output_schema={}, + ) + proc = train.input_processors["signal"] + before = copy.deepcopy(vars(proc)) + for split in ("validation", "test"): + held_out = create_sample_dataset( + samples=[{"patient_id": split, "signal": series([[1000]])}], + input_schema=schema, + output_schema={}, + input_processors=train.input_processors, + ) + result = held_out[0]["signal"] + result = result["value"] if isinstance(result, dict) else result + np.testing.assert_array_equal(result, [[98]]) + self.assertEqual(vars(proc), before) + + def test_sample_builder_save_load_preserves_normalization(self): + for alias in ("timeseries", "temporal_timeseries"): + with self.subTest(processor=alias), tempfile.TemporaryDirectory() as tmp: + builder = SampleBuilder( + {"signal": (alias, {"normalize_strategy": "standard"})}, {} + ) + builder.fit([{"signal": series([[10], [30]])}]) + path = str(Path(tmp) / "schema.pkl") + builder.save(path) + restored = SampleBuilder.load(path) + value = {"sample": pickle.dumps({"signal": series([[40]])})} + result = restored.transform(value)["signal"] + result = result["value"] if isinstance(result, dict) else result + np.testing.assert_array_equal(result, [[2]]) + + def test_cache_fingerprint_contains_training_statistics(self): + # BaseDataset.set_task serializes vars(processor) this way. + for cls in PROCESSORS: + with self.subTest(processor=cls.__name__): + fingerprints = [] + for values in ([[10], [30]], [[100], [300]], [[10], [30]]): + proc = cls(normalize_strategy="standard") + proc.fit([{"signal": series(values)}], "signal") + fingerprints.append( + json.dumps(vars(proc), sort_keys=True, default=str) + ) + self.assertNotEqual(fingerprints[0], fingerprints[1]) + self.assertEqual(fingerprints[0], fingerprints[2]) + + +if __name__ == "__main__": + unittest.main()