From d8e9878218a8d11a2d36fe1e3467f0a7f465c9be Mon Sep 17 00:00:00 2001 From: Kashif Rasul Date: Sat, 14 Dec 2019 10:34:38 +0100 Subject: [PATCH] moved transforms to their own module --- pts/__init__.py | 3 +- pts/dataset/__init__.py | 8 - pts/dataset/loader.py | 2 +- pts/dataset/stat.py | 2 +- pts/evaluation/backtest.py | 2 +- pts/feature/__init__.py | 27 +- pts/feature/transform.py | 960 ------------------ pts/model/deepar/deepar_estimator.py | 9 +- pts/model/estimator.py | 6 +- pts/model/predictor.py | 2 +- pts/transform/__init__.py | 59 ++ pts/transform/convert.py | 750 ++++++++++++++ .../dataset.py} | 7 +- pts/transform/feature.py | 264 +++++ pts/transform/field.py | 116 +++ pts/{dataset => transform}/sampler.py | 96 +- pts/transform/split.py | 548 ++++++++++ pts/transform/transform.py | 143 +++ setup.py | 1 + ...st_transformation.py => test_transform.py} | 402 ++++++-- 20 files changed, 2320 insertions(+), 1087 deletions(-) delete mode 100644 pts/feature/transform.py create mode 100644 pts/transform/__init__.py create mode 100644 pts/transform/convert.py rename pts/{dataset/transformed_dataset.py => transform/dataset.py} (90%) create mode 100644 pts/transform/feature.py create mode 100644 pts/transform/field.py rename pts/{dataset => transform}/sampler.py (55%) create mode 100644 pts/transform/split.py create mode 100644 pts/transform/transform.py rename test/{feature/test_transformation.py => test_transform.py} (59%) diff --git a/pts/__init__.py b/pts/__init__.py index 9ec9ed0..623c653 100644 --- a/pts/__init__.py +++ b/pts/__init__.py @@ -1 +1,2 @@ -from .trainer import Trainer \ No newline at end of file +from .trainer import Trainer +from .exception import assert_pts diff --git a/pts/dataset/__init__.py b/pts/dataset/__init__.py index e903add..10171b7 100644 --- a/pts/dataset/__init__.py +++ b/pts/dataset/__init__.py @@ -1,15 +1,7 @@ from .common import DataEntry, FieldName, Dataset from .list_dataset import ListDataset from .loader import TrainDataLoader, InferenceDataLoader -from .sampler import ( - InstanceSampler, - BucketInstanceSampler, - ExpectedNumInstanceSampler, - TestSplitSampler, - UniformSplitSampler, -) from .process import ProcessStartField, ProcessDataEntry from .utils import to_pandas from .stat import DatasetStatistics, ScaleHistogram, calculate_dataset_statistics from .artificial import constant_dataset -from .transformed_dataset import TransformedDataset \ No newline at end of file diff --git a/pts/dataset/loader.py b/pts/dataset/loader.py index 86b96bd..17e806e 100644 --- a/pts/dataset/loader.py +++ b/pts/dataset/loader.py @@ -8,8 +8,8 @@ import numpy as np import torch # First-party imports -from ..feature import Transformation from .common import DataEntry, Dataset +from pts.transform import Transformation DataBatch = Dict[str, Any] diff --git a/pts/dataset/stat.py b/pts/dataset/stat.py index ddcb665..9f81566 100644 --- a/pts/dataset/stat.py +++ b/pts/dataset/stat.py @@ -6,7 +6,7 @@ import numpy as np from tqdm import tqdm from .common import FieldName -from ..exception import assert_pts +from pts.exception import assert_pts class ScaleHistogram: diff --git a/pts/evaluation/backtest.py b/pts/evaluation/backtest.py index eae2af3..b9842ee 100644 --- a/pts/evaluation/backtest.py +++ b/pts/evaluation/backtest.py @@ -6,7 +6,7 @@ from typing import Dict, Iterator, NamedTuple, Optional, Tuple, Union import pandas as pd # First-party imports -from pts.feature import AdhocTransform +from pts.transform import AdhocTransform from pts.dataset import DataEntry, Dataset, TransformedDataset, InferenceDataLoader, DatasetStatistics, calculate_dataset_statistics from pts.model import Estimator, PTSEstimator, PTSPredictor, Predictor, Forecast from .evaluator import Evaluator diff --git a/pts/feature/__init__.py b/pts/feature/__init__.py index 91af972..7b46d61 100644 --- a/pts/feature/__init__.py +++ b/pts/feature/__init__.py @@ -10,30 +10,5 @@ from .time_feature import ( WeekOfYear, time_features_from_frequency_str, ) -from .transform import ( - AddAgeFeature, - AddConstFeature, - AddObservedValuesIndicator, - AddTimeFeatures, - AdhocTransform, - AsNumpyArray, - CanonicalInstanceSplitter, - Chain, - ConcatFeatures, - ExpandDimArray, - FilterTransformation, - FlatMapTransformation, - IdentityTransformation, - InstanceSplitter, - ListFeatures, - MapTransformation, - RemoveFields, - RenameFields, - SelectFields, - SetField, - SimpleTransformation, - SwapAxes, - Transformation, - VstackFeatures, -) + from .utils import get_granularity, get_seasonality \ No newline at end of file diff --git a/pts/feature/transform.py b/pts/feature/transform.py deleted file mode 100644 index 89da3d7..0000000 --- a/pts/feature/transform.py +++ /dev/null @@ -1,960 +0,0 @@ -from abc import ABC, abstractmethod -from collections import Counter -from functools import lru_cache, reduce -from typing import Any, Callable, Dict, Iterator, List, Optional, Tuple - -import numpy as np -import pandas as pd - -from ..exception import assert_pts -from ..dataset.sampler import InstanceSampler -from ..dataset import DataEntry -from .time_feature import TimeFeature - -MAX_IDLE_TRANSFORMS = 100 - - -@lru_cache(maxsize=10000) -def shift_timestamp(ts: pd.Timestamp, offset: int) -> pd.Timestamp: - try: - # this line looks innocent, but can create a date which is out of - # bounds. Values over year 9999 raise a ValueError and - # values over 2262-04-11 raise a pandas OutOfBoundsDatetime - return ts + offset * ts.freq - except (ValueError, pd._libs.OutOfBoundsDatetime) as ex: - raise Exception(ex) - - -def target_transformation_length( - target: np.array, pred_length: int, is_train: bool -) -> int: - return target.shape[-1] + (0 if is_train else pred_length) - - -class Transformation(ABC): - @abstractmethod - def __call__( - self, data_it: Iterator[DataEntry], is_train: bool - ) -> Iterator[DataEntry]: - pass - - def estimate(self, data_it: Iterator[DataEntry]) -> Iterator[DataEntry]: - return data_it # default is to pass through without estimation - - -class Chain(Transformation): - """ - Chain multiple transformations together. - """ - - def __init__(self, trans: List[Transformation]) -> None: - self.trans = trans - - def __call__( - self, data_it: Iterator[DataEntry], is_train: bool - ) -> Iterator[DataEntry]: - tmp = data_it - for t in self.trans: - tmp = t(tmp, is_train) - return tmp - - def estimate(self, data_it: Iterator[DataEntry]) -> Iterator[DataEntry]: - return reduce(lambda x, y: y.estimate(x), self.trans, data_it) - - -class IdentityTransformation(Transformation): - def __call__( - self, data_it: Iterator[DataEntry], is_train: bool - ) -> Iterator[DataEntry]: - return data_it - - -class MapTransformation(Transformation): - """ - Base class for Transformations that returns exactly one result per input in the stream. - """ - - def __call__(self, data_it: Iterator[DataEntry], is_train: bool) -> Iterator: - for data_entry in data_it: - try: - yield self.map_transform(data_entry.copy(), is_train) - except Exception as e: - raise e - - @abstractmethod - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - pass - - -class SimpleTransformation(MapTransformation): - """ - Element wise transformations that are the same in train and test mode - """ - - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - return self.transform(data) - - @abstractmethod - def transform(self, data: DataEntry) -> DataEntry: - pass - - -class AdhocTransform(SimpleTransformation): - """ - Applies a function as a transformation - This is called ad-hoc, because it is not serializable. - It is OK to use this for experiments and outside of a model pipeline that - needs to be serialized. - """ - - def __init__(self, func: Callable[[DataEntry], DataEntry]) -> None: - self.func = func - - def transform(self, data: DataEntry) -> DataEntry: - return self.func(data.copy()) - - -class FlatMapTransformation(Transformation): - """ - Transformations that yield zero or more results per input, but do not combine - elements from the input stream. - """ - - def __call__(self, data_it: Iterator[DataEntry], is_train: bool) -> Iterator: - num_idle_transforms = 0 - for data_entry in data_it: - num_idle_transforms += 1 - try: - for result in self.flatmap_transform(data_entry.copy(), is_train): - num_idle_transforms = 0 - yield result - except Exception as e: - raise e - if num_idle_transforms > MAX_IDLE_TRANSFORMS: - raise Exception( - f"Reached maximum number of idle transformation calls.\n" - f"This means the transformation looped over " - f"MAX_IDLE_TRANSFORMS={MAX_IDLE_TRANSFORMS} " - f"inputs without returning any output.\n" - f"This occurred in the following transformation:\n{self}" - ) - - @abstractmethod - def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: - pass - - -class FilterTransformation(FlatMapTransformation): - def __init__(self, condition: Callable[[DataEntry], bool]) -> None: - self.condition = condition - - def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: - if self.condition(data): - yield data - - -class RemoveFields(SimpleTransformation): - def __init__(self, field_names: List[str]) -> None: - self.field_names = field_names - - def transform(self, data: DataEntry) -> DataEntry: - for k in self.field_names: - if k in data.keys(): - del data[k] - return data - - -class SetField(SimpleTransformation): - """ - Sets a field in the dictionary with the given value. - Parameters - ---------- - output_field - Name of the field that will be set - value - Value to be set - """ - - def __init__(self, output_field: str, value: Any) -> None: - self.output_field = output_field - self.value = value - - def transform(self, data: DataEntry) -> DataEntry: - data[self.output_field] = self.value - return data - - -class SetFieldIfNotPresent(SimpleTransformation): - """ - Sets a field in the dictionary with the given value, in case it does not exist already - Parameters - ---------- - field - Name of the field that will be set - value - Value to be set - """ - - def __init__(self, field: str, value: Any) -> None: - self.output_field = field - self.value = value - - def transform(self, data: DataEntry) -> DataEntry: - if self.output_field not in data.keys(): - data[self.output_field] = self.value - return data - - -class AsNumpyArray(SimpleTransformation): - """ - Converts the value of a field into a numpy array. - Parameters - ---------- - expected_ndim - Expected number of dimensions. Throws an exception if the number of - dimensions does not match. - dtype - numpy dtype to use. - """ - - def __init__( - self, field: str, expected_ndim: int, dtype: np.dtype = np.float32 - ) -> None: - self.field = field - self.expected_ndim = expected_ndim - self.dtype = dtype - - def transform(self, data: DataEntry) -> DataEntry: - value = data[self.field] - if not isinstance(value, float): - # this lines produces "ValueError: setting an array element with a - # sequence" on our test - # value = np.asarray(value, dtype=np.float32) - # see https://stackoverflow.com/questions/43863748/ - value = np.asarray(list(value), dtype=self.dtype) - else: - # ugly: required as list conversion will fail in the case of a - # float - value = np.asarray(value, dtype=self.dtype) - assert_pts( - value.ndim >= self.expected_ndim, - 'Input for field "{self.field}" does not have the required' - "dimension (field: {self.field}, ndim observed: {value.ndim}, " - "expected ndim: {self.expected_ndim})", - value=value, - self=self, - ) - data[self.field] = value - return data - - -class ExpandDimArray(SimpleTransformation): - """ - Expand dims in the axis specified, if the axis is not present does nothing. - (This essentially calls np.expand_dims) - Parameters - ---------- - field - Field in dictionary to use - axis - Axis to expand (see np.expand_dims for details) - """ - - def __init__(self, field: str, axis: Optional[int] = None) -> None: - self.field = field - self.axis = axis - - def transform(self, data: DataEntry) -> DataEntry: - if self.axis is not None: - data[self.field] = np.expand_dims(data[self.field], axis=self.axis) - return data - - -class VstackFeatures(SimpleTransformation): - """ - Stack fields together using ``np.vstack``. - Fields with value ``None`` are ignored. - Parameters - ---------- - output_field - Field name to use for the output - input_fields - Fields to stack together - drop_inputs - If set to true the input fields will be dropped. - """ - - def __init__( - self, output_field: str, input_fields: List[str], drop_inputs: bool = True - ) -> None: - self.output_field = output_field - self.input_fields = input_fields - self.cols_to_drop = ( - [] - if not drop_inputs - else [fname for fname in self.input_fields if fname != output_field] - ) - - def transform(self, data: DataEntry) -> DataEntry: - r = [data[fname] for fname in self.input_fields if data[fname] is not None] - output = np.vstack(r) - data[self.output_field] = output - for fname in self.cols_to_drop: - del data[fname] - return data - - -class ConcatFeatures(SimpleTransformation): - """ - Concatenate fields together using ``np.concatenate``. - Fields with value ``None`` are ignored. - Parameters - ---------- - output_field - Field name to use for the output - input_fields - Fields to stack together - drop_inputs - If set to true the input fields will be dropped. - """ - - def __init__( - self, output_field: str, input_fields: List[str], drop_inputs: bool = True - ) -> None: - self.output_field = output_field - self.input_fields = input_fields - self.cols_to_drop = ( - [] - if not drop_inputs - else [fname for fname in self.input_fields if fname != output_field] - ) - - def transform(self, data: DataEntry) -> DataEntry: - r = [data[fname] for fname in self.input_fields if data[fname] is not None] - output = np.concatenate(r) - data[self.output_field] = output - for fname in self.cols_to_drop: - del data[fname] - return data - - -class SwapAxes(SimpleTransformation): - """ - Apply `np.swapaxes` to fields. - Parameters - ---------- - input_fields - Field to apply to - axes - Axes to use - """ - - def __init__(self, input_fields: List[str], axes: Tuple[int, int]) -> None: - self.input_fields = input_fields - self.axis1, self.axis2 = axes - - def transform(self, data: DataEntry) -> DataEntry: - for field in self.input_fields: - data[field] = self.swap(data[field]) - return data - - def swap(self, v): - if isinstance(v, np.ndarray): - return np.swapaxes(v, self.axis1, self.axis2) - if isinstance(v, list): - return [self.swap(x) for x in v] - else: - raise ValueError( - f"Unexpected field type {type(v).__name__}, expected " - f"np.ndarray or list[np.ndarray]" - ) - - -class ListFeatures(SimpleTransformation): - """ - Creates a new field which contains a list of features. - Parameters - ---------- - output_field - Field name for output - input_fields - Fields to combine into list - drop_inputs - If true the input fields will be removed from the result. - """ - - def __init__( - self, output_field: str, input_fields: List[str], drop_inputs: bool = True - ) -> None: - self.output_field = output_field - self.input_fields = input_fields - self.cols_to_drop = ( - [] - if not drop_inputs - else [fname for fname in self.input_fields if fname != output_field] - ) - - def transform(self, data: DataEntry) -> DataEntry: - data[self.output_field] = [data[fname] for fname in self.input_fields] - for fname in self.cols_to_drop: - del data[fname] - return data - - -class AddObservedValuesIndicator(SimpleTransformation): - """ - Replaces missing values in a numpy array (NaNs) with a dummy value and adds - an "observed"-indicator that is ``1`` when values are observed and ``0`` - when values are missing. - - Parameters - ---------- - target_field - Field for which missing values will be replaced - output_field - Field name to use for the indicator - dummy_value - Value to use for replacing missing values. - convert_nans - If set to true (default) missing values will be replaced. Otherwise - they will not be replaced. In any case the indicator is included in the - result. - dtype - """ - - def __init__( - self, - target_field: str, - output_field: str, - dummy_value: int = 0, - convert_nans: bool = True, - dtype: np.dtype = np.float32, - ) -> None: - self.dummy_value = dummy_value - self.target_field = target_field - self.output_field = output_field - self.convert_nans = convert_nans - self.dtype = dtype - - def transform(self, data: DataEntry) -> DataEntry: - value = data[self.target_field] - nan_indices = np.where(np.isnan(value)) - nan_entries = np.isnan(value) - - if self.convert_nans: - value[nan_indices] = self.dummy_value - - data[self.target_field] = value - # Invert bool array so that missing values are zeros and store as dtype - data[self.output_field] = np.invert(nan_entries).astype(self.dtype) - return data - - -class RenameFields(SimpleTransformation): - """ - Rename fields using a mapping - Parameters - ---------- - mapping - Name mapping `input_name -> output_name` - """ - - def __init__(self, mapping: Dict[str, str]) -> None: - self.mapping = mapping - values_count = Counter(mapping.values()) - for new_key, count in values_count.items(): - assert count == 1, f"Mapped key {new_key} occurs multiple time" - - def transform(self, data: DataEntry): - for key, new_key in self.mapping.items(): - if key not in data: - continue - assert new_key not in data - data[new_key] = data[key] - del data[key] - return data - - -class AddConstFeature(MapTransformation): - """ - Expands a `const` value along the time axis as a dynamic feature, where - the T-dimension is defined as the sum of the `pred_length` parameter and - the length of a time series specified by the `target_field`. - If `is_train=True` the feature matrix has the same length as the `target` field. - If `is_train=False` the feature matrix has length len(target) + pred_length - Parameters - ---------- - output_field - Field name for output. - target_field - Field containing the target array. The length of this array will be used. - pred_length - Prediction length (this is necessary since - features have to be available in the future) - const - Constant value to use. - dtype - Numpy dtype to use for resulting array. - """ - - def __init__( - self, - output_field: str, - target_field: str, - pred_length: int, - const: float = 1.0, - dtype: np.dtype = np.float32, - ) -> None: - self.pred_length = pred_length - self.const = const - self.dtype = dtype - self.output_field = output_field - self.target_field = target_field - - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - length = target_transformation_length( - data[self.target_field], self.pred_length, is_train=is_train - ) - data[self.output_field] = self.const * np.ones( - shape=(1, length), dtype=self.dtype - ) - return data - - -class AddTimeFeatures(MapTransformation): - """ - Adds a set of time features. - If `is_train=True` the feature matrix has the same length as the `target` field. - If `is_train=False` the feature matrix has length len(target) + pred_length - Parameters - ---------- - start_field - Field with the start time stamp of the time series - target_field - Field with the array containing the time series values - output_field - Field name for result. - time_features - list of time features to use. - pred_length - Prediction length - """ - - def __init__( - self, - start_field: str, - target_field: str, - output_field: str, - time_features: List[TimeFeature], - pred_length: int, - ) -> None: - self.date_features = time_features - self.pred_length = pred_length - self.start_field = start_field - self.target_field = target_field - self.output_field = output_field - self._min_time_point: Optional[pd.Timestamp] = None - self._max_time_point: Optional[pd.Timestamp] = None - self._full_range_date_features: Optional[np.ndarray] = None - self._date_index: Optional[pd.DatetimeIndex] = None - - def _update_cache(self, start: pd.Timestamp, length: int) -> None: - end = shift_timestamp(start, length) - if self._min_time_point is not None: - if self._min_time_point <= start and end <= self._max_time_point: - return - if self._min_time_point is None: - self._min_time_point = start - self._max_time_point = end - self._min_time_point = min(shift_timestamp(start, -50), self._min_time_point) - self._max_time_point = max(shift_timestamp(end, 50), self._max_time_point) - self.full_date_range = pd.date_range( - self._min_time_point, self._max_time_point, freq=start.freq - ) - self._full_range_date_features = ( - np.vstack( - [feat(self.full_date_range) for feat in self.date_features] - ) - if self.date_features - else None - ) - self._date_index = pd.Series( - index=self.full_date_range, data=np.arange(len(self.full_date_range)) - ) - - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - start = data[self.start_field] - length = target_transformation_length( - data[self.target_field], self.pred_length, is_train=is_train - ) - self._update_cache(start, length) - i0 = self._date_index[start] - features = ( - self._full_range_date_features[..., i0 : i0 + length] - if self.date_features - else None - ) - data[self.output_field] = features - return data - - -class AddAgeFeature(MapTransformation): - """ - Adds an 'age' feature to the data_entry. - The age feature starts with a small value at the start of the time series - and grows over time. - If `is_train=True` the age feature has the same length as the `target` field. - If `is_train=False` the age feature has length len(target) + pred_length - Parameters - ---------- - target_field - Field with target values (array) of time series - output_field - Field name to use for the output. - pred_length - Prediction length - log_scale - If set to true the age feature grows logarithmically otherwise linearly over time. - dtype - """ - - def __init__( - self, - target_field: str, - output_field: str, - pred_length: int, - log_scale: bool = True, - dtype: np.dtype = np.float32, - ) -> None: - self.pred_length = pred_length - self.target_field = target_field - self.feature_name = output_field - self.log_scale = log_scale - self._age_feature = np.zeros(0) - self.dtype = dtype - - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - length = target_transformation_length( - data[self.target_field], self.pred_length, is_train=is_train - ) - - if self.log_scale: - age = np.log10(2.0 + np.arange(length, dtype=self.dtype)) - else: - age = np.arange(length, dtype=self.dtype) - - data[self.feature_name] = age.reshape((1, length)) - - return data - - -class InstanceSplitter(FlatMapTransformation): - """ - Selects training instances, by slicing the target and other time series - like arrays at random points in training mode or at the last time point in - prediction mode. Assumption is that all time like arrays start at the same - time point. - The target and each time_series_field is removed and instead two - corresponding fields with prefix `past_` and `future_` are included. E.g. - If the target array is one-dimensional, the resulting instance has shape - (len_target). In the multi-dimensional case, the instance has shape (dim, - len_target). - target -> past_target and future_target - The transformation also adds a field 'past_is_pad' that indicates whether - values where padded or not. - Convention: time axis is always the last axis. - - Parameters - ---------- - target_field - field containing the target - is_pad_field - output field indicating whether padding happened - start_field - field containing the start date of the time series - forecast_start_field - output field that will contain the time point where the forecast starts - train_sampler - instance sampler that provides sampling indices given a time-series - past_length - length of the target seen before making prediction - future_length - length of the target that must be predicted - batch_first - whether to have time series output in (time, dimension) or in - (dimension, time) layout - time_series_fields - fields that contains time-series, they are split in the same interval - as the target - pick_incomplete - whether training examples can be sampled with only a part of - past_length time-units - present for the time series. This is useful to train models for - cold-start. In such case, is_pad_out contains an indicator whether - data is padded or not. - """ - - def __init__( - self, - target_field: str, - is_pad_field: str, - start_field: str, - forecast_start_field: str, - train_sampler: InstanceSampler, - past_length: int, - future_length: int, - batch_first: bool = True, - time_series_fields: Optional[List[str]] = None, - pick_incomplete: bool = True, - ) -> None: - - assert future_length > 0 - - self.train_sampler = train_sampler - self.past_length = past_length - self.future_length = future_length - self.batch_first = batch_first - self.ts_fields = time_series_fields if time_series_fields is not None else [] - self.target_field = target_field - self.is_pad_field = is_pad_field - self.start_field = start_field - self.forecast_start_field = forecast_start_field - self.pick_incomplete = pick_incomplete - - def _past(self, col_name): - return f"past_{col_name}" - - def _future(self, col_name): - return f"future_{col_name}" - - def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: - pl = self.future_length - slice_cols = self.ts_fields + [self.target_field] - target = data[self.target_field] - - len_target = target.shape[-1] - - if is_train: - if len_target < self.future_length: - # We currently cannot handle time series that are shorter than - # the prediction length during training, so we just skip these. - # If we want to include them we would need to pad and to mask - # the loss. - sampling_indices: List[int] = [] - else: - if self.pick_incomplete: - sampling_indices = self.train_sampler( - target, 0, len_target - self.future_length - ) - else: - sampling_indices = self.train_sampler( - target, self.past_length, len_target - self.future_length - ) - else: - sampling_indices = [len_target] - for i in sampling_indices: - pad_length = max(self.past_length - i, 0) - if not self.pick_incomplete: - assert pad_length == 0 - d = data.copy() - for ts_field in slice_cols: - if i > self.past_length: - # truncate to past_length - past_piece = d[ts_field][..., i - self.past_length : i] - elif i < self.past_length: - pad_block = np.zeros( - d[ts_field].shape[:-1] + (pad_length,), dtype=d[ts_field].dtype - ) - past_piece = np.concatenate( - [pad_block, d[ts_field][..., :i]], axis=-1 - ) - else: - past_piece = d[ts_field][..., :i] - d[self._past(ts_field)] = past_piece - d[self._future(ts_field)] = d[ts_field][..., i : i + pl] - del d[ts_field] - pad_indicator = np.zeros(self.past_length) - if pad_length > 0: - pad_indicator[:pad_length] = 1 - - if self.batch_first: - for ts_field in slice_cols: - d[self._past(ts_field)] = d[self._past(ts_field)].transpose() - d[self._future(ts_field)] = d[self._future(ts_field)].transpose() - - d[self._past(self.is_pad_field)] = pad_indicator - d[self.forecast_start_field] = shift_timestamp(d[self.start_field], i) - yield d - - -class CanonicalInstanceSplitter(FlatMapTransformation): - """ - Selects instances, by slicing the target and other time series - like arrays at random points in training mode or at the last time point in - prediction mode. Assumption is that all time like arrays start at the same - time point. - In training mode, the returned instances contain past_`target_field` - as well as past_`time_series_fields`. - In prediction mode, one can set `use_prediction_features` to get - future_`time_series_fields`. - If the target array is one-dimensional, the `target_field` in the resulting instance has shape - (`instance_length`). In the multi-dimensional case, the instance has shape (`dim`, `instance_length`), - where `dim` can also take a value of 1. - In the case of insufficient number of time series values, the - transformation also adds a field 'past_is_pad' that indicates whether - values where padded or not, and the value is padded with - `default_pad_value` with a default value 0. - This is done only if `allow_target_padding` is `True`, - and the length of `target` is smaller than `instance_length`. - Parameters - ---------- - target_field - fields that contains time-series - is_pad_field - output field indicating whether padding happened - start_field - field containing the start date of the time series - forecast_start_field - field containing the forecast start date - instance_sampler - instance sampler that provides sampling indices given a time-series - instance_length - length of the target seen before making prediction - batch_first - whether to have time series output in (time, dimension) or in - (dimension, time) layout - time_series_fields - fields that contains time-series, they are split in the same interval - as the target - allow_target_padding - flag to allow padding - pad_value - value to be used for padding - use_prediction_features - flag to indicate if prediction range features should be returned - prediction_length - length of the prediction range, must be set if - use_prediction_features is True - """ - - def __init__( - self, - target_field: str, - is_pad_field: str, - start_field: str, - forecast_start_field: str, - instance_sampler: InstanceSampler, - instance_length: int, - batch_first: bool = True, - time_series_fields: List[str] = [], - allow_target_padding: bool = False, - pad_value: float = 0.0, - use_prediction_features: bool = False, - prediction_length: Optional[int] = None, - ) -> None: - self.instance_sampler = instance_sampler - self.instance_length = instance_length - self.batch_first = batch_first - self.dynamic_feature_fields = time_series_fields - self.target_field = target_field - self.allow_target_padding = allow_target_padding - self.pad_value = pad_value - self.is_pad_field = is_pad_field - self.start_field = start_field - self.forecast_start_field = forecast_start_field - - assert ( - not use_prediction_features or prediction_length is not None - ), "You must specify `prediction_length` if `use_prediction_features`" - - self.use_prediction_features = use_prediction_features - self.prediction_length = prediction_length - - def _past(self, col_name): - return f"past_{col_name}" - - def _future(self, col_name): - return f"future_{col_name}" - - def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: - ts_fields = self.dynamic_feature_fields + [self.target_field] - ts_target = data[self.target_field] - - len_target = ts_target.shape[-1] - - if is_train: - if len_target < self.instance_length: - sampling_indices = ( - # Returning [] for all time series will cause this to be in loop forever! - [len_target] - if self.allow_target_padding - else [] - ) - else: - sampling_indices = self.instance_sampler( - ts_target, self.instance_length, len_target - ) - else: - sampling_indices = [len_target] - - for i in sampling_indices: - d = data.copy() - - pad_length = max(self.instance_length - i, 0) - - # update start field - d[self.start_field] = shift_timestamp( - data[self.start_field], i - self.instance_length - ) - - # set is_pad field - is_pad = np.zeros(self.instance_length) - if pad_length > 0: - is_pad[:pad_length] = 1 - d[self.is_pad_field] = is_pad - - # update time series fields - for ts_field in ts_fields: - full_ts = data[ts_field] - if pad_length > 0: - pad_pre = self.pad_value * np.ones( - shape=full_ts.shape[:-1] + (pad_length,) - ) - past_ts = np.concatenate([pad_pre, full_ts[..., :i]], axis=-1) - else: - past_ts = full_ts[..., (i - self.instance_length) : i] - - past_ts = past_ts.transpose() if self.batch_first else past_ts - d[self._past(ts_field)] = past_ts - - if self.use_prediction_features and not is_train: - if not ts_field == self.target_field: - future_ts = full_ts[..., i : i + self.prediction_length] - future_ts = ( - future_ts.transpose() if self.batch_first else future_ts - ) - d[self._future(ts_field)] = future_ts - - del d[ts_field] - - d[self.forecast_start_field] = shift_timestamp( - d[self.start_field], self.instance_length - ) - - yield d - - -class SelectFields(MapTransformation): - """ - Only keep the listed fields - Parameters - ---------- - input_fields - List of fields to keep. - """ - - def __init__(self, input_fields: List[str]) -> None: - self.input_fields = input_fields - - def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: - return {f: data[f] for f in self.input_fields} diff --git a/pts/model/deepar/deepar_estimator.py b/pts/model/deepar/deepar_estimator.py index 3aaa898..81ab8eb 100644 --- a/pts/model/deepar/deepar_estimator.py +++ b/pts/model/deepar/deepar_estimator.py @@ -9,7 +9,8 @@ from pts import Trainer from pts.feature import ( TimeFeature, get_lags_for_frequency, - time_features_from_frequency_str, + time_features_from_frequency_str) +from pts.transform import ( Transformation, Chain, RemoveFields, @@ -20,13 +21,13 @@ from pts.feature import ( AddAgeFeature, VstackFeatures, InstanceSplitter, + ExpectedNumInstanceSampler, ) -from pts.dataset import FieldName, ExpectedNumInstanceSampler -from pts.model import PTSEstimator, Predictor, PTSPredictor +from pts.dataset import FieldName +from pts.model import PTSEstimator, Predictor, PTSPredictor, copy_parameters from pts.modules import DistributionOutput, StudentTOutput from .deepar_network import DeepARTrainingNetwork, DeepARPredictionNetwork -from ..utils import copy_parameters class DeepAREstimator(PTSEstimator): def __init__( diff --git a/pts/model/estimator.py b/pts/model/estimator.py index 7edbaf1..99e1665 100644 --- a/pts/model/estimator.py +++ b/pts/model/estimator.py @@ -5,11 +5,11 @@ import numpy as np import torch import torch.nn as nn -from ..dataset import Dataset, TrainDataLoader -from ..feature import Transformation +from pts.dataset import Dataset, TrainDataLoader +from pts.transform import Transformation +from pts import Trainer from .predictor import Predictor -from ..trainer import Trainer from .utils import get_module_forward_input_names diff --git a/pts/model/predictor.py b/pts/model/predictor.py index 7561eae..79b2399 100644 --- a/pts/model/predictor.py +++ b/pts/model/predictor.py @@ -6,7 +6,7 @@ import torch import torch.nn as nn from pts.dataset import Dataset, DataEntry, InferenceDataLoader -from pts.feature import Transformation +from pts.transform import Transformation from .forecast import Forecast from .forecast_generator import ForecastGenerator, SampleForecastGenerator diff --git a/pts/transform/__init__.py b/pts/transform/__init__.py new file mode 100644 index 0000000..4d1f56b --- /dev/null +++ b/pts/transform/__init__.py @@ -0,0 +1,59 @@ +from .transform import ( + Transformation, + Chain, + Identity, + MapTransformation, + SimpleTransformation, + AdhocTransform, + FlatMapTransformation, + FilterTransformation, +) + +from .convert import ( + AsNumpyArray, + ExpandDimArray, + VstackFeatures, + ConcatFeatures, + SwapAxes, + ListFeatures, + TargetDimIndicator, + SampleTargetDim, + CDFtoGaussianTransform, + cdf_to_gaussian_forward_transform, +) + +from .dataset import TransformedDataset + +from .feature import ( + target_transformation_length, + AddObservedValuesIndicator, + AddConstFeature, + AddTimeFeatures, + AddAgeFeature, +) + +from .field import ( + RemoveFields, + RenameFields, + SetField, + SetFieldIfNotPresent, + SelectFields, +) + + +from .sampler import ( + InstanceSampler, + UniformSplitSampler, + TestSplitSampler, + ExpectedNumInstanceSampler, + BucketInstanceSampler, + ContinuousTimePointSampler, + ContinuousTimeUniformSampler, +) + +from .split import ( + shift_timestamp, + InstanceSplitter, + CanonicalInstanceSplitter, + ContinuousTimeInstanceSplitter, +) diff --git a/pts/transform/convert.py b/pts/transform/convert.py new file mode 100644 index 0000000..a95cbc4 --- /dev/null +++ b/pts/transform/convert.py @@ -0,0 +1,750 @@ +# Copyright 2018 Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# A copy of the License is located at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# or in the "license" file accompanying this file. This file is distributed +# on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either +# express or implied. See the License for the specific language governing +# permissions and limitations under the License. + +from typing import Iterator, List, Tuple, Optional + +import numpy as np +from scipy.special import erf, erfinv + +import torch + +from pts.exception import assert_pts +from pts.dataset import DataEntry + +from pts.dataset import DataEntry + +from .transform import ( + SimpleTransformation, + MapTransformation, + FlatMapTransformation, +) + + +class AsNumpyArray(SimpleTransformation): + """ + Converts the value of a field into a numpy array. + + Parameters + ---------- + expected_ndim + Expected number of dimensions. Throws an exception if the number of + dimensions does not match. + dtype + numpy dtype to use. + """ + + def __init__( + self, field: str, expected_ndim: int, dtype: np.dtype = np.float32 + ) -> None: + self.field = field + self.expected_ndim = expected_ndim + self.dtype = dtype + + def transform(self, data: DataEntry) -> DataEntry: + value = data[self.field] + if not isinstance(value, float): + # this lines produces "ValueError: setting an array element with a + # sequence" on our test + # value = np.asarray(value, dtype=np.float32) + # see https://stackoverflow.com/questions/43863748/ + value = np.asarray(list(value), dtype=self.dtype) + else: + # ugly: required as list conversion will fail in the case of a + # float + value = np.asarray(value, dtype=self.dtype) + assert_pts( + value.ndim >= self.expected_ndim, + 'Input for field "{self.field}" does not have the required' + "dimension (field: {self.field}, ndim observed: {value.ndim}, " + "expected ndim: {self.expected_ndim})", + value=value, + self=self, + ) + data[self.field] = value + return data + + +class ExpandDimArray(SimpleTransformation): + """ + Expand dims in the axis specified, if the axis is not present does nothing. + (This essentially calls np.expand_dims) + + Parameters + ---------- + field + Field in dictionary to use + axis + Axis to expand (see np.expand_dims for details) + """ + + def __init__(self, field: str, axis: Optional[int] = None) -> None: + self.field = field + self.axis = axis + + def transform(self, data: DataEntry) -> DataEntry: + if self.axis is not None: + data[self.field] = np.expand_dims(data[self.field], axis=self.axis) + return data + + +class VstackFeatures(SimpleTransformation): + """ + Stack fields together using ``np.vstack``. + + Fields with value ``None`` are ignored. + + Parameters + ---------- + output_field + Field name to use for the output + input_fields + Fields to stack together + drop_inputs + If set to true the input fields will be dropped. + """ + + def __init__( + self, + output_field: str, + input_fields: List[str], + drop_inputs: bool = True, + ) -> None: + self.output_field = output_field + self.input_fields = input_fields + self.cols_to_drop = ( + [] + if not drop_inputs + else [ + fname for fname in self.input_fields if fname != output_field + ] + ) + + def transform(self, data: DataEntry) -> DataEntry: + r = [ + data[fname] + for fname in self.input_fields + if data[fname] is not None + ] + output = np.vstack(r) + data[self.output_field] = output + for fname in self.cols_to_drop: + del data[fname] + return data + + +class ConcatFeatures(SimpleTransformation): + """ + Concatenate fields together using ``np.concatenate``. + + Fields with value ``None`` are ignored. + + Parameters + ---------- + output_field + Field name to use for the output + input_fields + Fields to stack together + drop_inputs + If set to true the input fields will be dropped. + """ + + def __init__( + self, + output_field: str, + input_fields: List[str], + drop_inputs: bool = True, + ) -> None: + self.output_field = output_field + self.input_fields = input_fields + self.cols_to_drop = ( + [] + if not drop_inputs + else [ + fname for fname in self.input_fields if fname != output_field + ] + ) + + def transform(self, data: DataEntry) -> DataEntry: + r = [ + data[fname] + for fname in self.input_fields + if data[fname] is not None + ] + output = np.concatenate(r) + data[self.output_field] = output + for fname in self.cols_to_drop: + del data[fname] + return data + + +class SwapAxes(SimpleTransformation): + """ + Apply `np.swapaxes` to fields. + + Parameters + ---------- + input_fields + Field to apply to + axes + Axes to use + """ + + def __init__(self, input_fields: List[str], axes: Tuple[int, int]) -> None: + self.input_fields = input_fields + self.axis1, self.axis2 = axes + + def transform(self, data: DataEntry) -> DataEntry: + for field in self.input_fields: + data[field] = self.swap(data[field]) + return data + + def swap(self, v): + if isinstance(v, np.ndarray): + return np.swapaxes(v, self.axis1, self.axis2) + if isinstance(v, list): + return [self.swap(x) for x in v] + else: + raise ValueError( + f"Unexpected field type {type(v).__name__}, expected " + f"np.ndarray or list[np.ndarray]" + ) + + +class ListFeatures(SimpleTransformation): + """ + Creates a new field which contains a list of features. + + Parameters + ---------- + output_field + Field name for output + input_fields + Fields to combine into list + drop_inputs + If true the input fields will be removed from the result. + """ + + def __init__( + self, + output_field: str, + input_fields: List[str], + drop_inputs: bool = True, + ) -> None: + self.output_field = output_field + self.input_fields = input_fields + self.cols_to_drop = ( + [] + if not drop_inputs + else [ + fname for fname in self.input_fields if fname != output_field + ] + ) + + def transform(self, data: DataEntry) -> DataEntry: + data[self.output_field] = [data[fname] for fname in self.input_fields] + for fname in self.cols_to_drop: + del data[fname] + return data + + +class TargetDimIndicator(SimpleTransformation): + """ + Label-encoding of the target dimensions. + """ + + def __init__(self, field_name: str, target_field: str) -> None: + self.field_name = field_name + self.target_field = target_field + + def transform(self, data: DataEntry) -> DataEntry: + data[self.field_name] = np.arange(0, data[self.target_field].shape[0]) + return data + + +class SampleTargetDim(FlatMapTransformation): + """ + Samples random dimensions from the target at training time. + """ + + def __init__( + self, + field_name: str, + target_field: str, + observed_values_field: str, + num_samples: int, + shuffle: bool = True, + ) -> None: + self.field_name = field_name + self.target_field = target_field + self.observed_values_field = observed_values_field + self.num_samples = num_samples + self.shuffle = shuffle + + def flatmap_transform( + self, data: DataEntry, is_train: bool, slice_future_target: bool = True + ) -> Iterator[DataEntry]: + if not is_train: + yield data + else: + # (target_dim,) + target_dimensions = data[self.field_name] + + if self.shuffle: + np.random.shuffle(target_dimensions) + + target_dimensions = target_dimensions[: self.num_samples] + + data[self.field_name] = target_dimensions + # (seq_len, target_dim) -> (seq_len, num_samples) + + for field in [ + f"past_{self.target_field}", + f"future_{self.target_field}", + f"past_{self.observed_values_field}", + f"future_{self.observed_values_field}", + ]: + data[field] = data[field][:, target_dimensions] + + yield data + + +class CDFtoGaussianTransform(MapTransformation): + """ + Marginal transformation that transforms the target via an empirical CDF + to a standard gaussian as described here: https://arxiv.org/abs/1910.03002 + + To be used in conjunction with a multivariate gaussian to from a copula. + Note that this transformation is currently intended for multivariate + targets only. + """ + + def __init__( + self, + target_dim: int, + target_field: str, + observed_values_field: str, + cdf_suffix="_cdf", + max_context_length: Optional[int] = None, + ) -> None: + """ + Constructor for CDFtoGaussianTransform. + + Parameters + ---------- + target_dim + Dimensionality of the target. + target_field + Field that will be transformed. + observed_values_field + Field that indicates observed values. + cdf_suffix + Suffix to mark the field with the transformed target. + max_context_length + Sets the maximum context length for the empirical CDF. + """ + self.target_field = target_field + self.past_target_field = "past_" + self.target_field + self.future_target_field = "future_" + self.target_field + self.past_observed_field = f"past_{observed_values_field}" + self.sort_target_field = f"past_{target_field}_sorted" + self.slopes_field = "slopes" + self.intercepts_field = "intercepts" + self.cdf_suffix = cdf_suffix + self.max_context_length = max_context_length + self.target_dim = target_dim + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + self._preprocess_data(data, is_train=is_train) + self._calc_pw_linear_params(data) + + for target_field in [self.past_target_field, self.future_target_field]: + data[target_field + self.cdf_suffix] = self.standard_gaussian_ppf( + self._empirical_cdf_forward_transform( + data[self.sort_target_field], + data[target_field], + data[self.slopes_field], + data[self.intercepts_field], + ) + ) + return data + + def _preprocess_data(self, data: DataEntry, is_train: bool): + """ + Performs several preprocess operations for computing the empirical CDF. + 1) Reshaping the data. + 2) Normalizing the target length. + 3) Adding noise to avoid zero slopes (training only) + 4) Sorting the target to compute the empirical CDF + + Parameters + ---------- + data + DataEntry with input data. + is_train + if is_train is True, this function adds noise to the target to + avoid zero slopes in the piece-wise linear function. + Returns + ------- + + """ + # (target_length, target_dim) + past_target_vec = data[self.past_target_field].copy() + + # pick only observed values + target_length, target_dim = past_target_vec.shape + + # (target_length, target_dim) + past_observed = (data[self.past_observed_field] > 0) * ( + data["past_is_pad"].reshape((-1, 1)) == 0 + ) + assert past_observed.ndim == 2 + assert target_dim == self.target_dim + + past_target_vec = past_target_vec[past_observed.min(axis=1)] + + assert past_target_vec.ndim == 2 + assert past_target_vec.shape[1] == self.target_dim + + expected_length = ( + target_length + if self.max_context_length is None + else self.max_context_length + ) + + if target_length != expected_length: + # Fills values in the case where past_target_vec.shape[-1] < + # target_length + # as dataset.loader.BatchBuffer does not support varying shapes + past_target_vec = CDFtoGaussianTransform._fill( + past_target_vec, expected_length + ) + + # sorts along the time dimension to compute empirical CDF of each + # dimension + if is_train: + past_target_vec = self._add_noise(past_target_vec) + + past_target_vec.sort(axis=0) + + assert past_target_vec.shape == (expected_length, self.target_dim) + + data[self.sort_target_field] = past_target_vec + + def _calc_pw_linear_params(self, data: DataEntry): + """ + Calculates the piece-wise linear parameters to interpolate between + the observed values in the empirical CDF. + + Once current limitation is that we use a zero slope line as the last + piece. Thus, we cannot forecast anything higher than the highest + observed value. + + Parameters + ---------- + data + Input data entry containing a sorted target field. + + Returns + ------- + + """ + sorted_target = data[self.sort_target_field] + sorted_target_length, target_dim = sorted_target.shape + + quantiles = np.stack( + [np.arange(sorted_target_length) for _ in range(target_dim)], + axis=1, + ) / float(sorted_target_length) + + x_diff = np.diff(sorted_target, axis=0) + y_diff = np.diff(quantiles, axis=0) + + # Calculate slopes of the pw-linear pieces. + slopes = np.where( + x_diff == 0.0, np.zeros_like(x_diff), y_diff / x_diff + ) + + zeroes = np.zeros_like(np.expand_dims(slopes[0, :], axis=0)) + slopes = np.append(slopes, zeroes, axis=0) + + # Calculate intercepts of the pw-linear pieces. + intercepts = quantiles - slopes * sorted_target + + # Populate new fields with the piece-wise linear parameters. + data[self.slopes_field] = slopes + data[self.intercepts_field] = intercepts + + def _empirical_cdf_forward_transform( + self, + sorted_values: np.ndarray, + values: np.ndarray, + slopes: np.ndarray, + intercepts: np.ndarray, + ) -> np.ndarray: + """ + Applies the empirical CDF forward transformation. + + Parameters + ---------- + sorted_values + Sorted target vector. + values + Values (real valued) that will be transformed to empirical CDF + values. + slopes + Slopes of the piece-wise linear function. + intercepts + Intercepts of the piece-wise linear function. + + Returns + ------- + quantiles + Empirical CDF quantiles in [0, 1] interval with winzorized cutoff. + + """ + m = sorted_values.shape[0] + quantiles = self._forward_transform( + sorted_values, values, slopes, intercepts + ) + + quantiles = np.clip( + quantiles, self.winsorized_cutoff(m), 1 - self.winsorized_cutoff(m) + ) + return quantiles + + @staticmethod + def _add_noise(x: np.array) -> np.array: + scale_noise = 0.2 + std = np.sqrt( + (np.square(x - x.mean(axis=1, keepdims=True))).mean( + axis=1, keepdims=True + ) + ) + noise = np.random.normal( + loc=np.zeros_like(x), scale=np.ones_like(x) * std * scale_noise + ) + x = x + noise + return x + + @staticmethod + def _search_sorted( + sorted_vec: np.array, to_insert_vec: np.array + ) -> np.array: + """ + Finds the indices of the active piece-wise linear function. + + Parameters + ---------- + sorted_vec + Sorted target vector. + to_insert_vec + Vector for which the indicies of the active linear functions + will be computed + + Returns + ------- + indices + Indices mapping to the active linear function. + """ + indices_left = np.searchsorted(sorted_vec, to_insert_vec, side="left") + indices_right = np.searchsorted( + sorted_vec, to_insert_vec, side="right" + ) + + indices = indices_left + (indices_right - indices_left) // 2 + indices = indices - 1 + indices = np.minimum(indices, len(sorted_vec) - 1) + indices[indices < 0] = 0 + return indices + + def _forward_transform( + self, + sorted_vec: np.array, + target: np.array, + slopes: np.array, + intercepts: np.array, + ) -> np.array: + """ + Applies the forward transformation to the marginals of the multivariate + target. Target (real valued) -> empirical cdf [0, 1] + + Parameters + ---------- + sorted_vec + Sorted (past) target vector. + target + Target that will be transformed. + slopes + Slopes of the piece-wise linear function. + intercepts + Intercepts of the piece-wise linear function + + Returns + ------- + transformed_target + Transformed target vector. + """ + transformed = list() + for sorted, t, slope, intercept in zip( + sorted_vec.transpose(), + target.transpose(), + slopes.transpose(), + intercepts.transpose(), + ): + indices = self._search_sorted(sorted, t) + transformed_value = slope[indices] * t + intercept[indices] + transformed.append(transformed_value) + return np.array(transformed).transpose() + + @staticmethod + def standard_gaussian_cdf(x: np.array) -> np.array: + u = x / (np.sqrt(2.0)) + return (erf(u) + 1.0) / 2.0 + + @staticmethod + def standard_gaussian_ppf(y: np.array) -> np.array: + y_clipped = np.clip(y, a_min=1.0e-6, a_max=1.0 - 1.0e-6) + return np.sqrt(2.0) * erfinv(2.0 * y_clipped - 1.0) + + @staticmethod + def winsorized_cutoff(m: np.array) -> np.array: + """ + Apply truncation to the empirical CDF estimator to reduce variance as + described here: https://arxiv.org/abs/0903.0649 + + Parameters + ---------- + m + Input array with empirical CDF values. + + Returns + ------- + res + Truncated empirical CDf values. + """ + res = 1 / (4 * m ** 0.25 * np.sqrt(3.14 * np.log(m))) + assert 0 < res < 1 + return res + + @staticmethod + def _fill(target: np.ndarray, expected_length: int) -> np.ndarray: + """ + Makes sure target has at least expected_length time-units by repeating + it or using zeros. + + Parameters + ---------- + target : shape (seq_len, dim) + expected_length + + Returns + ------- + array of shape (target_length, dim) + """ + + current_length, target_dim = target.shape + if current_length == 0: + # todo handle the case with no observation better, + # we could use dataset statistics but for now we use zeros + filled_target = np.zeros((expected_length, target_dim)) + elif current_length < expected_length: + filled_target = np.vstack( + [target for _ in range(expected_length // current_length + 1)] + ) + filled_target = filled_target[:expected_length] + elif current_length > expected_length: + filled_target = target[-expected_length:] + else: + filled_target = target + + assert filled_target.shape == (expected_length, target_dim) + + return filled_target + + +def cdf_to_gaussian_forward_transform( + input_batch: DataEntry, outputs: torch.Tensor +) -> np.ndarray: + """ + Forward transformation of the CDFtoGaussianTransform. + + Parameters + ---------- + input_batch + Input data to the predictor. + outputs + Predictor outputs. + Returns + ------- + outputs + Forward transformed outputs. + + """ + + def _empirical_cdf_inverse_transform( + batch_target_sorted: torch.Tensor, + batch_predictions: torch.Tensor, + slopes: torch.Tensor, + intercepts: torch.Tensor, + ) -> np.ndarray: + """ + Apply forward transformation of the empirical CDF. + + Parameters + ---------- + batch_target_sorted + Sorted targets of the input batch. + batch_predictions + Predictions of the underlying probability distribution + slopes + Slopes of the piece-wise linear function. + intercepts + Intercepts of the piece-wise linear function. + + Returns + ------- + outputs + Forward transformed outputs. + + """ + slopes = slopes.cpu().numpy() + intercepts = intercepts.cpu().numpy() + + batch_target_sorted = batch_target_sorted.cpu().numpy() + _, num_timesteps, _ = batch_target_sorted.shape + indices = np.floor(batch_predictions * num_timesteps) + # indices = indices - 1 + # for now project into [0, 1] + indices = np.clip(indices, 0, num_timesteps - 1) + indices = indices.astype(np.int) + + transformed = np.where( + np.take_along_axis(slopes, indices, axis=1) != 0.0, + (batch_predictions - np.take_along_axis(intercepts, indices, axis=1)) + / np.take_along_axis(slopes, indices, axis=1), + np.take_along_axis(batch_target_sorted, indices, axis=1), + ) + return transformed + + # applies inverse cdf to all outputs + _, samples, _, _ = outputs.shape + for sample_index in range(0, samples): + outputs[:, sample_index, :, :] = _empirical_cdf_inverse_transform( + input_batch["past_target_sorted"], + CDFtoGaussianTransform.standard_gaussian_cdf( + outputs[:, sample_index, :, :] + ), + input_batch["slopes"], + input_batch["intercepts"], + ) + return outputs diff --git a/pts/dataset/transformed_dataset.py b/pts/transform/dataset.py similarity index 90% rename from pts/dataset/transformed_dataset.py rename to pts/transform/dataset.py index cf0538b..98d0e9e 100644 --- a/pts/dataset/transformed_dataset.py +++ b/pts/transform/dataset.py @@ -1,8 +1,7 @@ from typing import Iterator, List -from pts.feature import Chain, Transformation - -from .common import DataEntry, Dataset +from pts.dataset import DataEntry, Dataset +from .transform import Chain, Transformation class TransformedDataset(Dataset): @@ -11,6 +10,8 @@ class TransformedDataset(Dataset): element in the base_dataset. This only supports SimpleTransformations, which do the same thing at prediction and training time. + + Parameters ---------- base_dataset diff --git a/pts/transform/feature.py b/pts/transform/feature.py new file mode 100644 index 0000000..4defd5d --- /dev/null +++ b/pts/transform/feature.py @@ -0,0 +1,264 @@ +# Copyright 2018 Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# A copy of the License is located at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# or in the "license" file accompanying this file. This file is distributed +# on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either +# express or implied. See the License for the specific language governing +# permissions and limitations under the License. + +from typing import List + +import numpy as np +import pandas as pd + +from pts.dataset import DataEntry +from pts.feature import TimeFeature + +from .transform import SimpleTransformation, MapTransformation +from .split import shift_timestamp + + +def target_transformation_length( + target: np.array, pred_length: int, is_train: bool +) -> int: + return target.shape[-1] + (0 if is_train else pred_length) + + +class AddObservedValuesIndicator(SimpleTransformation): + """ + Replaces missing values in a numpy array (NaNs) with a dummy value and adds + an "observed"-indicator that is ``1`` when values are observed and ``0`` + when values are missing. + + + Parameters + ---------- + target_field + Field for which missing values will be replaced + output_field + Field name to use for the indicator + dummy_value + Value to use for replacing missing values. + convert_nans + If set to true (default) missing values will be replaced. Otherwise + they will not be replaced. In any case the indicator is included in the + result. + """ + + def __init__( + self, + target_field: str, + output_field: str, + dummy_value: int = 0, + convert_nans: bool = True, + dtype: np.dtype = np.float32, + ) -> None: + self.dummy_value = dummy_value + self.target_field = target_field + self.output_field = output_field + self.convert_nans = convert_nans + self.dtype = dtype + + def transform(self, data: DataEntry) -> DataEntry: + value = data[self.target_field] + nan_indices = np.where(np.isnan(value)) + nan_entries = np.isnan(value) + + if self.convert_nans: + value[nan_indices] = self.dummy_value + + data[self.target_field] = value + # Invert bool array so that missing values are zeros and store as float + data[self.output_field] = np.invert(nan_entries).astype(self.dtype) + return data + + +class AddConstFeature(MapTransformation): + """ + Expands a `const` value along the time axis as a dynamic feature, where + the T-dimension is defined as the sum of the `pred_length` parameter and + the length of a time series specified by the `target_field`. + + If `is_train=True` the feature matrix has the same length as the `target` field. + If `is_train=False` the feature matrix has length len(target) + pred_length + + Parameters + ---------- + output_field + Field name for output. + target_field + Field containing the target array. The length of this array will be used. + pred_length + Prediction length (this is necessary since + features have to be available in the future) + const + Constant value to use. + dtype + Numpy dtype to use for resulting array. + """ + + def __init__( + self, + output_field: str, + target_field: str, + pred_length: int, + const: float = 1.0, + dtype: np.dtype = np.float32, + ) -> None: + self.pred_length = pred_length + self.const = const + self.dtype = dtype + self.output_field = output_field + self.target_field = target_field + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + length = target_transformation_length( + data[self.target_field], self.pred_length, is_train=is_train + ) + data[self.output_field] = self.const * np.ones( + shape=(1, length), dtype=self.dtype + ) + return data + + +class AddTimeFeatures(MapTransformation): + """ + Adds a set of time features. + + If `is_train=True` the feature matrix has the same length as the `target` field. + If `is_train=False` the feature matrix has length len(target) + pred_length + + Parameters + ---------- + start_field + Field with the start time stamp of the time series + target_field + Field with the array containing the time series values + output_field + Field name for result. + time_features + list of time features to use. + pred_length + Prediction length + """ + + def __init__( + self, + start_field: str, + target_field: str, + output_field: str, + time_features: List[TimeFeature], + pred_length: int, + ) -> None: + self.date_features = time_features + self.pred_length = pred_length + self.start_field = start_field + self.target_field = target_field + self.output_field = output_field + self._min_time_point: pd.Timestamp = None + self._max_time_point: pd.Timestamp = None + self._full_range_date_features: np.ndarray = None + self._date_index: pd.DatetimeIndex = None + + def _update_cache(self, start: pd.Timestamp, length: int) -> None: + end = shift_timestamp(start, length) + if self._min_time_point is not None: + if self._min_time_point <= start and end <= self._max_time_point: + return + if self._min_time_point is None: + self._min_time_point = start + self._max_time_point = end + self._min_time_point = min( + shift_timestamp(start, -50), self._min_time_point + ) + self._max_time_point = max( + shift_timestamp(end, 50), self._max_time_point + ) + self.full_date_range = pd.date_range( + self._min_time_point, self._max_time_point, freq=start.freq + ) + self._full_range_date_features = ( + np.vstack( + [feat(self.full_date_range) for feat in self.date_features] + ) + if self.date_features + else None + ) + self._date_index = pd.Series( + index=self.full_date_range, + data=np.arange(len(self.full_date_range)), + ) + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + start = data[self.start_field] + length = target_transformation_length( + data[self.target_field], self.pred_length, is_train=is_train + ) + self._update_cache(start, length) + i0 = self._date_index[start] + features = ( + self._full_range_date_features[..., i0 : i0 + length] + if self.date_features + else None + ) + data[self.output_field] = features + return data + + +class AddAgeFeature(MapTransformation): + """ + Adds an 'age' feature to the data_entry. + + The age feature starts with a small value at the start of the time series + and grows over time. + + If `is_train=True` the age feature has the same length as the `target` + field. + If `is_train=False` the age feature has length len(target) + pred_length + + Parameters + ---------- + target_field + Field with target values (array) of time series + output_field + Field name to use for the output. + pred_length + Prediction length + log_scale + If set to true the age feature grows logarithmically otherwise linearly + over time. + """ + + def __init__( + self, + target_field: str, + output_field: str, + pred_length: int, + log_scale: bool = True, + dtype: np.dtype = np.float32, + ) -> None: + self.pred_length = pred_length + self.target_field = target_field + self.feature_name = output_field + self.log_scale = log_scale + self._age_feature = np.zeros(0) + self.dtype = dtype + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + length = target_transformation_length( + data[self.target_field], self.pred_length, is_train=is_train + ) + + if self.log_scale: + age = np.log10(2.0 + np.arange(length, dtype=self.dtype)) + else: + age = np.arange(length, dtype=self.dtype) + + data[self.feature_name] = age.reshape((1, length)) + + return data diff --git a/pts/transform/field.py b/pts/transform/field.py new file mode 100644 index 0000000..06a36ba --- /dev/null +++ b/pts/transform/field.py @@ -0,0 +1,116 @@ +# Copyright 2018 Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# A copy of the License is located at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# or in the "license" file accompanying this file. This file is distributed +# on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either +# express or implied. See the License for the specific language governing +# permissions and limitations under the License. + +from collections import Counter +from typing import Any, Dict, List + +from pts.dataset import DataEntry + +from .transform import SimpleTransformation, MapTransformation + + +class RenameFields(SimpleTransformation): + """ + Rename fields using a mapping + + Parameters + ---------- + mapping + Name mapping `input_name -> output_name` + """ + + def __init__(self, mapping: Dict[str, str]) -> None: + self.mapping = mapping + values_count = Counter(mapping.values()) + for new_key, count in values_count.items(): + assert count == 1, f"Mapped key {new_key} occurs multiple time" + + def transform(self, data: DataEntry): + for key, new_key in self.mapping.items(): + if key not in data: + continue + assert new_key not in data + data[new_key] = data[key] + del data[key] + return data + + +class RemoveFields(SimpleTransformation): + def __init__(self, field_names: List[str]) -> None: + self.field_names = field_names + + def transform(self, data: DataEntry) -> DataEntry: + for k in self.field_names: + if k in data.keys(): + del data[k] + return data + + +class SetField(SimpleTransformation): + """ + Sets a field in the dictionary with the given value. + + Parameters + ---------- + output_field + Name of the field that will be set + value + Value to be set + """ + + def __init__(self, output_field: str, value: Any) -> None: + self.output_field = output_field + self.value = value + + def transform(self, data: DataEntry) -> DataEntry: + data[self.output_field] = self.value + return data + + +class SetFieldIfNotPresent(SimpleTransformation): + """Sets a field in the dictionary with the given value, in case it does not + exist already. + + Parameters + ---------- + output_field + Name of the field that will be set + value + Value to be set + """ + + def __init__(self, field: str, value: Any) -> None: + self.output_field = field + self.value = value + + def transform(self, data: DataEntry) -> DataEntry: + if self.output_field not in data.keys(): + data[self.output_field] = self.value + return data + + +class SelectFields(MapTransformation): + """ + Only keep the listed fields + + Parameters + ---------- + input_fields + List of fields to keep. + """ + + def __init__(self, input_fields: List[str]) -> None: + self.input_fields = input_fields + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + return {f: data[f] for f in self.input_fields} diff --git a/pts/dataset/sampler.py b/pts/transform/sampler.py similarity index 55% rename from pts/dataset/sampler.py rename to pts/transform/sampler.py index 18460a8..a46b6c1 100644 --- a/pts/dataset/sampler.py +++ b/pts/transform/sampler.py @@ -1,20 +1,54 @@ +# Copyright 2018 Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# A copy of the License is located at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# or in the "license" file accompanying this file. This file is distributed +# on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either +# express or implied. See the License for the specific language governing +# permissions and limitations under the License. + from abc import ABC, abstractmethod -from typing import List, Union import numpy as np -from .stat import ScaleHistogram +from pts.dataset.stat import ScaleHistogram -class InstanceSampler(ABC): - @abstractmethod - def __call__(self, ts: np.ndarray, a: int, b: int) -> Union[np.ndarray, List[int]]: - pass +class InstanceSampler: + """ + An InstanceSampler is called with the time series and the valid + index bounds a, b and should return a set of indices a <= i <= b + at which training instances will be generated. + + The object should be called with: + + Parameters + ---------- + ts + target that should be sampled with shape (dim, seq_len) + a + first index of the target that can be sampled + b + last index of the target that can be sampled + + Returns + ------- + np.ndarray + Selected points to sample + """ + + def __call__(self, ts: np.ndarray, a: int, b: int) -> np.ndarray: + raise NotImplementedError() class UniformSplitSampler(InstanceSampler): """ Samples each point with the same fixed probability. + Parameters ---------- p @@ -26,7 +60,9 @@ class UniformSplitSampler(InstanceSampler): self.lookup = np.arange(2 ** 13) def __call__(self, ts: np.ndarray, a: int, b: int) -> np.ndarray: - assert a <= b + assert ( + a <= b + ), "First index must be less than or equal to the last index." while ts.shape[-1] >= len(self.lookup): self.lookup = np.arange(2 * len(self.lookup)) mask = np.random.uniform(low=0.0, high=1.0, size=b - a + 1) < self.p @@ -39,6 +75,9 @@ class TestSplitSampler(InstanceSampler): splitting i.e. the forecast point for the time series. """ + def __init__(self) -> None: + pass + def __call__(self, ts: np.ndarray, a: int, b: int) -> np.ndarray: return np.array([b]) @@ -48,8 +87,10 @@ class ExpectedNumInstanceSampler(InstanceSampler): Keeps track of the average time series length and adjusts the probability per time point such that on average `num_instances` training examples are generated per time series. + Parameters ---------- + num_instances number of training examples generated per time series on average """ @@ -78,7 +119,9 @@ class BucketInstanceSampler(InstanceSampler): This sample can be used when working with a set of time series that have a skewed distributions. For instance, if the dataset contains many time series with small values and few with large values. + The probability of sampling from bucket i is the inverse of its number of elements. + Parameters ---------- scale_histogram @@ -92,10 +135,47 @@ class BucketInstanceSampler(InstanceSampler): self.scale_histogram = scale_histogram self.lookup = np.arange(2 ** 13) - def __call__(self, ts: np.ndarray, a: int, b: int) -> np.ndarray: + def __call__(self, ts: np.ndarray, a: int, b: int) -> None: while ts.shape[-1] >= len(self.lookup): self.lookup = np.arange(2 * len(self.lookup)) p = 1.0 / self.scale_histogram.count(ts) mask = np.random.uniform(low=0.0, high=1.0, size=b - a + 1) < p indices = self.lookup[a : a + len(mask)][mask] return indices + + +class ContinuousTimePointSampler(ABC): + """ + Abstract class for "continuous time" samplers, which, given a lower bound + and upper bound, sample "points" (events) in continuous time from a + specified interval. + """ + + def __init__(self, num_instances: int) -> None: + self.num_instances = num_instances + + @abstractmethod + def __call__(self, a: float, b: float) -> np.ndarray: + """ + Returns random points in the real interval between :code:`a` and + :code:`b`. + + Parameters + ---------- + a + The lower bound (minimum time value that a sampled point can take) + b + Upper bound. Must be greater than a. + """ + pass + + +class ContinuousTimeUniformSampler(ContinuousTimePointSampler): + """ + Implements a simple random sampler to sample points in the continuous + interval between :code:`a` and :code:`b`. + """ + + def __call__(self, a: float, b: float) -> np.ndarray: + assert a <= b, "Interval start time must be before interval end time." + return np.random.rand(self.num_instances) * (b - a) + a diff --git a/pts/transform/split.py b/pts/transform/split.py new file mode 100644 index 0000000..38bf86c --- /dev/null +++ b/pts/transform/split.py @@ -0,0 +1,548 @@ +# Copyright 2018 Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# A copy of the License is located at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# or in the "license" file accompanying this file. This file is distributed +# on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either +# express or implied. See the License for the specific language governing +# permissions and limitations under the License. + +from functools import lru_cache +from typing import Iterator, List, Optional + +import numpy as np +import pandas as pd + +from pts.dataset import DataEntry, FieldName + +from .transform import FlatMapTransformation +from .sampler import InstanceSampler, ContinuousTimePointSampler + + +def shift_timestamp(ts: pd.Timestamp, offset: int) -> pd.Timestamp: + """ + Computes a shifted timestamp. + + Basic wrapping around pandas ``ts + offset`` with caching and exception + handling. + """ + return _shift_timestamp_helper(ts, ts.freq, offset) + + +@lru_cache(maxsize=10000) +def _shift_timestamp_helper( + ts: pd.Timestamp, freq: str, offset: int +) -> pd.Timestamp: + """ + We are using this helper function which explicitly uses the frequency as a + parameter, because the frequency is not included in the hash of a time + stamp. + + I.e. + pd.Timestamp(x, freq='1D') and pd.Timestamp(x, freq='1min') + + hash to the same value. + """ + try: + # this line looks innocent, but can create a date which is out of + # bounds values over year 9999 raise a ValueError + # values over 2262-04-11 raise a pandas OutOfBoundsDatetime + return ts + offset * ts.freq + except (ValueError, pd._libs.OutOfBoundsDatetime) as ex: + raise Exception(ex) + + +class InstanceSplitter(FlatMapTransformation): + """ + Selects training instances, by slicing the target and other time series + like arrays at random points in training mode or at the last time point in + prediction mode. Assumption is that all time like arrays start at the same + time point. + + The target and each time_series_field is removed and instead two + corresponding fields with prefix `past_` and `future_` are included. E.g. + + If the target array is one-dimensional, the resulting instance has shape + (len_target). In the multi-dimensional case, the instance has shape (dim, + len_target). + + target -> past_target and future_target + + The transformation also adds a field 'past_is_pad' that indicates whether + values where padded or not. + + Convention: time axis is always the last axis. + + Parameters + ---------- + + target_field + field containing the target + is_pad_field + output field indicating whether padding happened + start_field + field containing the start date of the time series + forecast_start_field + output field that will contain the time point where the forecast starts + train_sampler + instance sampler that provides sampling indices given a time-series + past_length + length of the target seen before making prediction + future_length + length of the target that must be predicted + output_NTC + whether to have time series output in (time, dimension) or in + (dimension, time) layout + time_series_fields + fields that contains time-series, they are split in the same interval + as the target + pick_incomplete + whether training examples can be sampled with only a part of + past_length time-units + present for the time series. This is useful to train models for + cold-start. In such case, is_pad_out contains an indicator whether + data is padded or not. + """ + + def __init__( + self, + target_field: str, + is_pad_field: str, + start_field: str, + forecast_start_field: str, + train_sampler: InstanceSampler, + past_length: int, + future_length: int, + output_NTC: bool = True, + time_series_fields: Optional[List[str]] = None, + pick_incomplete: bool = True, + ) -> None: + + assert future_length > 0 + + self.train_sampler = train_sampler + self.past_length = past_length + self.future_length = future_length + self.output_NTC = output_NTC + self.ts_fields = ( + time_series_fields if time_series_fields is not None else [] + ) + self.target_field = target_field + self.is_pad_field = is_pad_field + self.start_field = start_field + self.forecast_start_field = forecast_start_field + self.pick_incomplete = pick_incomplete + + def _past(self, col_name): + return f"past_{col_name}" + + def _future(self, col_name): + return f"future_{col_name}" + + def flatmap_transform( + self, data: DataEntry, is_train: bool + ) -> Iterator[DataEntry]: + pl = self.future_length + slice_cols = self.ts_fields + [self.target_field] + target = data[self.target_field] + + len_target = target.shape[-1] + + if is_train: + if len_target < self.future_length: + # We currently cannot handle time series that are shorter than + # the prediction length during training, so we just skip these. + # If we want to include them we would need to pad and to mask + # the loss. + sampling_indices: List[int] = [] + else: + if self.pick_incomplete: + sampling_indices = self.train_sampler( + target, 0, len_target - self.future_length + ) + else: + sampling_indices = self.train_sampler( + target, + self.past_length, + len_target - self.future_length, + ) + else: + sampling_indices = [len_target] + for i in sampling_indices: + pad_length = max(self.past_length - i, 0) + if not self.pick_incomplete: + assert pad_length == 0 + d = data.copy() + for ts_field in slice_cols: + if i > self.past_length: + # truncate to past_length + past_piece = d[ts_field][..., i - self.past_length : i] + elif i < self.past_length: + pad_block = np.zeros( + d[ts_field].shape[:-1] + (pad_length,), + dtype=d[ts_field].dtype, + ) + past_piece = np.concatenate( + [pad_block, d[ts_field][..., :i]], axis=-1 + ) + else: + past_piece = d[ts_field][..., :i] + d[self._past(ts_field)] = past_piece + d[self._future(ts_field)] = d[ts_field][..., i : i + pl] + del d[ts_field] + pad_indicator = np.zeros(self.past_length) + if pad_length > 0: + pad_indicator[:pad_length] = 1 + + if self.output_NTC: + for ts_field in slice_cols: + d[self._past(ts_field)] = d[ + self._past(ts_field) + ].transpose() + d[self._future(ts_field)] = d[ + self._future(ts_field) + ].transpose() + + d[self._past(self.is_pad_field)] = pad_indicator + d[self.forecast_start_field] = shift_timestamp( + d[self.start_field], i + ) + yield d + + +class CanonicalInstanceSplitter(FlatMapTransformation): + """ + Selects instances, by slicing the target and other time series + like arrays at random points in training mode or at the last time point in + prediction mode. Assumption is that all time like arrays start at the same + time point. + + In training mode, the returned instances contain past_`target_field` + as well as past_`time_series_fields`. + + In prediction mode, one can set `use_prediction_features` to get + future_`time_series_fields`. + + If the target array is one-dimensional, the `target_field` in the resulting instance has shape + (`instance_length`). In the multi-dimensional case, the instance has shape (`dim`, `instance_length`), + where `dim` can also take a value of 1. + + In the case of insufficient number of time series values, the + transformation also adds a field 'past_is_pad' that indicates whether + values where padded or not, and the value is padded with + `default_pad_value` with a default value 0. + This is done only if `allow_target_padding` is `True`, + and the length of `target` is smaller than `instance_length`. + + Parameters + ---------- + target_field + fields that contains time-series + is_pad_field + output field indicating whether padding happened + start_field + field containing the start date of the time series + forecast_start_field + field containing the forecast start date + instance_sampler + instance sampler that provides sampling indices given a time-series + instance_length + length of the target seen before making prediction + output_NTC + whether to have time series output in (time, dimension) or in + (dimension, time) layout + time_series_fields + fields that contains time-series, they are split in the same interval + as the target + allow_target_padding + flag to allow padding + pad_value + value to be used for padding + use_prediction_features + flag to indicate if prediction range features should be returned + prediction_length + length of the prediction range, must be set if + use_prediction_features is True + """ + + def __init__( + self, + target_field: str, + is_pad_field: str, + start_field: str, + forecast_start_field: str, + instance_sampler: InstanceSampler, + instance_length: int, + output_NTC: bool = True, + time_series_fields: List[str] = [], + allow_target_padding: bool = False, + pad_value: float = 0.0, + use_prediction_features: bool = False, + prediction_length: Optional[int] = None, + ) -> None: + self.instance_sampler = instance_sampler + self.instance_length = instance_length + self.output_NTC = output_NTC + self.dynamic_feature_fields = time_series_fields + self.target_field = target_field + self.allow_target_padding = allow_target_padding + self.pad_value = pad_value + self.is_pad_field = is_pad_field + self.start_field = start_field + self.forecast_start_field = forecast_start_field + + assert ( + not use_prediction_features or prediction_length is not None + ), "You must specify `prediction_length` if `use_prediction_features`" + + self.use_prediction_features = use_prediction_features + self.prediction_length = prediction_length + + def _past(self, col_name): + return f"past_{col_name}" + + def _future(self, col_name): + return f"future_{col_name}" + + def flatmap_transform( + self, data: DataEntry, is_train: bool + ) -> Iterator[DataEntry]: + ts_fields = self.dynamic_feature_fields + [self.target_field] + ts_target = data[self.target_field] + + len_target = ts_target.shape[-1] + + if is_train: + if len_target < self.instance_length: + sampling_indices = ( + # Returning [] for all time series will cause this to be in loop forever! + [len_target] + if self.allow_target_padding + else [] + ) + else: + sampling_indices = self.instance_sampler( + ts_target, self.instance_length, len_target + ) + else: + sampling_indices = [len_target] + + for i in sampling_indices: + d = data.copy() + + pad_length = max(self.instance_length - i, 0) + + # update start field + d[self.start_field] = shift_timestamp( + data[self.start_field], i - self.instance_length + ) + + # set is_pad field + is_pad = np.zeros(self.instance_length) + if pad_length > 0: + is_pad[:pad_length] = 1 + d[self.is_pad_field] = is_pad + + # update time series fields + for ts_field in ts_fields: + full_ts = data[ts_field] + if pad_length > 0: + pad_pre = self.pad_value * np.ones( + shape=full_ts.shape[:-1] + (pad_length,) + ) + past_ts = np.concatenate( + [pad_pre, full_ts[..., :i]], axis=-1 + ) + else: + past_ts = full_ts[..., (i - self.instance_length) : i] + + past_ts = past_ts.transpose() if self.output_NTC else past_ts + d[self._past(ts_field)] = past_ts + + if self.use_prediction_features and not is_train: + if not ts_field == self.target_field: + future_ts = full_ts[ + ..., i : i + self.prediction_length + ] + future_ts = ( + future_ts.transpose() + if self.output_NTC + else future_ts + ) + d[self._future(ts_field)] = future_ts + + del d[ts_field] + + d[self.forecast_start_field] = shift_timestamp( + d[self.start_field], self.instance_length + ) + + yield d + + +class ContinuousTimeInstanceSplitter(FlatMapTransformation): + """ + Selects training instances by slicing "intervals" from a continous-time + process instantiation. Concretely, the input data is expected to describe an + instantiation from a point (or jump) process, with the "target" + identifying inter-arrival times and other features (marks), as described + in detail below. + + The splitter will then take random points in continuous time from each + given observation, and return a (variable-length) array of points in + the past (context) and the future (prediction) intervals. + + The transformation is analogous to its discrete counterpart + `InstanceSplitter` except that + + - It does not allow "incomplete" records. That is, the past and future + intervals sampled are always complete + - Outputs a (T, C) layout. + - Does not accept `time_series_fields` (i.e., only accepts target fields) as these + would typically not be available in TPP data. + + The target arrays are expected to have (2, T) layout where the first axis + corresponds to the (i) interarrival times between consecutive points, in + order and (ii) integer identifiers of marks (from {0, 1, ..., :code:`num_marks`}). + The returned arrays will have (T, 2) layout. + + For example, the array below corresponds to a target array where points with timestamps + 0.5, 1.1, and 1.5 were observed belonging to categories (marks) 3, 1 and 0 + respectively: :code:`[[0.5, 0.6, 0.4], [3, 1, 0]]`. + + Parameters + ---------- + past_interval_length + length of the interval seen before making prediction + future_interval_length + length of the interval that must be predicted + train_sampler + instance sampler that provides sampling indices given a time-series + target_field + field containing the target + start_field + field containing the start date of the of the point process observation + end_field + field containing the end date of the point process observation + forecast_start_field + output field that will contain the time point where the forecast starts + """ + + def __init__( + self, + past_interval_length: float, + future_interval_length: float, + train_sampler: ContinuousTimePointSampler, + target_field: str = FieldName.TARGET, + start_field: str = FieldName.START, + end_field: str = "end", + forecast_start_field: str = FieldName.FORECAST_START, + ) -> None: + + assert ( + future_interval_length > 0 + ), "Prediction interval must have length greater than 0." + + self.train_sampler = train_sampler + self.past_interval_length = past_interval_length + self.future_interval_length = future_interval_length + self.target_field = target_field + self.start_field = start_field + self.end_field = end_field + self.forecast_start_field = forecast_start_field + + # noinspection PyMethodMayBeStatic + def _mask_sorted(self, a: np.ndarray, lb: float, ub: float): + start = np.searchsorted(a, lb) + end = np.searchsorted(a, ub) + return np.arange(start, end) + + def flatmap_transform( + self, data: DataEntry, is_train: bool + ) -> Iterator[DataEntry]: + + assert data[self.start_field].freq == data[self.end_field].freq + + total_interval_length = ( + data[self.end_field] - data[self.start_field] + ) / data[self.start_field].freq.delta + + # sample forecast start times in continuous time + if is_train: + if total_interval_length < ( + self.future_interval_length + self.past_interval_length + ): + sampling_times: np.ndarray = np.array([]) + else: + sampling_times = self.train_sampler( + self.past_interval_length, + total_interval_length - self.future_interval_length, + ) + else: + sampling_times = np.array([total_interval_length]) + + ia_times = data[self.target_field][0, :] + marks = data[self.target_field][1:, :] + + ts = np.cumsum(ia_times) + assert ts[-1] < total_interval_length, ( + "Target interarrival times provided are inconsistent with " + "start and end timestamps." + ) + + # select field names that will be included in outputs + keep_cols = { + k: v + for k, v in data.items() + if k not in [self.target_field, self.start_field, self.end_field] + } + + for future_start in sampling_times: + + r: DataEntry = dict() + + past_start = future_start - self.past_interval_length + future_end = future_start + self.future_interval_length + + assert past_start >= 0 + + past_mask = self._mask_sorted(ts, past_start, future_start) + + past_ia_times = np.diff(np.r_[0, ts[past_mask] - past_start])[ + np.newaxis + ] + + r[f"past_{self.target_field}"] = np.concatenate( + [past_ia_times, marks[:, past_mask]], axis=0 + ).transpose() + + r["past_valid_length"] = np.array([len(past_mask)]) + + r[self.forecast_start_field] = ( + data[self.start_field] + + data[self.start_field].freq.delta * future_start + ) + + if is_train: # include the future only if is_train + assert future_end <= total_interval_length + + future_mask = self._mask_sorted(ts, future_start, future_end) + + future_ia_times = np.diff( + np.r_[0, ts[future_mask] - future_start] + )[np.newaxis] + + r[f"future_{self.target_field}"] = np.concatenate( + [future_ia_times, marks[:, future_mask]], axis=0 + ).transpose() + + r["future_valid_length"] = np.array([len(future_mask)]) + + # include other fields + r.update(keep_cols.copy()) + + yield r diff --git a/pts/transform/transform.py b/pts/transform/transform.py new file mode 100644 index 0000000..b1d061d --- /dev/null +++ b/pts/transform/transform.py @@ -0,0 +1,143 @@ +from abc import ABC, abstractmethod +from typing import Callable, Iterator, List +from functools import reduce + + +from pts.dataset import DataEntry + + +MAX_IDLE_TRANSFORMS = 100 + + +class Transformation(ABC): + @abstractmethod + def __call__( + self, data_it: Iterator[DataEntry], is_train: bool + ) -> Iterator[DataEntry]: + pass + + def estimate(self, data_it: Iterator[DataEntry]) -> Iterator[DataEntry]: + return data_it # default is to pass through without estimation + + def chain(self, other: "Transformation") -> "Chain": + return Chain(self, other) + + def __add__(self, other: "Transformation") -> "Chain": + return self.chain(other) + + +class Chain(Transformation): + """ + Chain multiple transformations together. + """ + + def __init__(self, trans: List[Transformation]) -> None: + self.transformations = [] + for transformation in trans: + # flatten chains + if isinstance(transformation, Chain): + self.transformations.extend(transformation.transformations) + else: + self.transformations.append(transformation) + + def __call__( + self, data_it: Iterator[DataEntry], is_train: bool + ) -> Iterator[DataEntry]: + tmp = data_it + for t in self.transformations: + tmp = t(tmp, is_train) + return tmp + + def estimate(self, data_it: Iterator[DataEntry]) -> Iterator[DataEntry]: + return reduce(lambda x, y: y.estimate(x), self.transformations, data_it) + + +class Identity(Transformation): + def __call__( + self, data_it: Iterator[DataEntry], is_train: bool + ) -> Iterator[DataEntry]: + return data_it + + +class MapTransformation(Transformation): + """ + Base class for Transformations that returns exactly one result per input in the stream. + """ + + def __call__(self, data_it: Iterator[DataEntry], is_train: bool) -> Iterator: + for data_entry in data_it: + try: + yield self.map_transform(data_entry.copy(), is_train) + except Exception as e: + raise e + + @abstractmethod + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + pass + + +class SimpleTransformation(MapTransformation): + """ + Element wise transformations that are the same in train and test mode + """ + + def map_transform(self, data: DataEntry, is_train: bool) -> DataEntry: + return self.transform(data) + + @abstractmethod + def transform(self, data: DataEntry) -> DataEntry: + pass + + +class AdhocTransform(SimpleTransformation): + """ + Applies a function as a transformation + This is called ad-hoc, because it is not serializable. + It is OK to use this for experiments and outside of a model pipeline that + needs to be serialized. + """ + + def __init__(self, func: Callable[[DataEntry], DataEntry]) -> None: + self.func = func + + def transform(self, data: DataEntry) -> DataEntry: + return self.func(data.copy()) + + +class FlatMapTransformation(Transformation): + """ + Transformations that yield zero or more results per input, but do not combine + elements from the input stream. + """ + + def __call__(self, data_it: Iterator[DataEntry], is_train: bool) -> Iterator: + num_idle_transforms = 0 + for data_entry in data_it: + num_idle_transforms += 1 + try: + for result in self.flatmap_transform(data_entry.copy(), is_train): + num_idle_transforms = 0 + yield result + except Exception as e: + raise e + if num_idle_transforms > MAX_IDLE_TRANSFORMS: + raise Exception( + f"Reached maximum number of idle transformation calls.\n" + f"This means the transformation looped over " + f"MAX_IDLE_TRANSFORMS={MAX_IDLE_TRANSFORMS} " + f"inputs without returning any output.\n" + f"This occurred in the following transformation:\n{self}" + ) + + @abstractmethod + def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: + pass + + +class FilterTransformation(FlatMapTransformation): + def __init__(self, condition: Callable[[DataEntry], bool]) -> None: + self.condition = condition + + def flatmap_transform(self, data: DataEntry, is_train: bool) -> Iterator[DataEntry]: + if self.condition(data): + yield data diff --git a/setup.py b/setup.py index 69b7b9e..cb2a6aa 100644 --- a/setup.py +++ b/setup.py @@ -19,6 +19,7 @@ setup( 'holidays', 'numpy', 'pandas', + 'scipy', 'tqdm', 'ujson', 'pydantic', diff --git a/test/feature/test_transformation.py b/test/test_transform.py similarity index 59% rename from test/feature/test_transformation.py rename to test/test_transform.py index 6b1840f..d231c07 100644 --- a/test/feature/test_transformation.py +++ b/test/test_transform.py @@ -5,29 +5,18 @@ import pytest import numpy as np import pandas as pd +import torch + from pts.dataset import ( ProcessStartField, FieldName, ListDataset, - ExpectedNumInstanceSampler, - UniformSplitSampler, - ScaleHistogram, - BucketInstanceSampler, + DataEntry, calculate_dataset_statistics, + ScaleHistogram, ) -from pts.feature import ( - Chain, - AddTimeFeatures, - DayOfWeek, - DayOfMonth, - MonthOfYear, - AddAgeFeature, - VstackFeatures, - InstanceSplitter, - CanonicalInstanceSplitter, - AddObservedValuesIndicator, -) - +from pts import transform +from pts.feature import time_feature FREQ = "1D" @@ -72,16 +61,14 @@ def test_align_timestamp(): @pytest.mark.parametrize("start", TEST_VALUES["start"]) def test_AddTimeFeatures(start, target, is_train): pred_length = 13 - t = AddTimeFeatures( + t = transform.AddTimeFeatures( start_field=FieldName.START, target_field=FieldName.TARGET, output_field="myout", pred_length=pred_length, - time_features=[DayOfWeek(), DayOfMonth()], + time_features=[time_feature.DayOfWeek(), time_feature.DayOfMonth()], ) - # TODO - # assert_serializable(t) data = {"start": start, "target": target} res = t.map_transform(data, is_train=is_train) @@ -89,8 +76,8 @@ def test_AddTimeFeatures(start, target, is_train): expected_length = len(target) + (0 if is_train else pred_length) assert mat.shape == (2, expected_length) tmp_idx = pd.date_range(start=start, freq=start.freq, periods=expected_length) - assert np.alltrue(mat[0] == DayOfWeek()(tmp_idx)) - assert np.alltrue(mat[1] == DayOfMonth()(tmp_idx)) + assert np.alltrue(mat[0] == time_feature.DayOfWeek()(tmp_idx)) + assert np.alltrue(mat[1] == time_feature.DayOfMonth()(tmp_idx)) @pytest.mark.parametrize("is_train", TEST_VALUES["is_train"]) @@ -98,7 +85,7 @@ def test_AddTimeFeatures(start, target, is_train): @pytest.mark.parametrize("start", TEST_VALUES["start"]) def test_AddTimeFeatures_empty_time_features(start, target, is_train): pred_length = 13 - t = AddTimeFeatures( + t = transform.AddTimeFeatures( start_field=FieldName.START, target_field=FieldName.TARGET, output_field="myout", @@ -106,8 +93,6 @@ def test_AddTimeFeatures_empty_time_features(start, target, is_train): time_features=[], ) - # assert_serializable(t) - data = {"start": start, "target": target} res = t.map_transform(data, is_train=is_train) assert res["myout"] is None @@ -118,14 +103,13 @@ def test_AddTimeFeatures_empty_time_features(start, target, is_train): @pytest.mark.parametrize("start", TEST_VALUES["start"]) def test_AddAgeFeatures(start, target, is_train): pred_length = 13 - t = AddAgeFeature( + t = transform.AddAgeFeature( pred_length=pred_length, target_field=FieldName.TARGET, output_field="age", log_scale=True, ) - # assert_serializable(t) data = {"start": start, "target": target} out = t.map_transform(data, is_train=is_train) @@ -143,19 +127,18 @@ def test_AddAgeFeatures(start, target, is_train): def test_InstanceSplitter(start, target, is_train): train_length = 100 pred_length = 13 - t = InstanceSplitter( + t = transform.InstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - train_sampler=UniformSplitSampler(p=1.0), + train_sampler=transform.UniformSplitSampler(p=1.0), past_length=train_length, future_length=pred_length, time_series_fields=["some_time_feature"], pick_incomplete=True, ) - # assert_serializable(t) other_feat = np.arange(len(target) + 100) data = { @@ -204,12 +187,12 @@ def test_CanonicalInstanceSplitter( ): train_length = 100 pred_length = 13 - t = CanonicalInstanceSplitter( + t = transform.CanonicalInstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - instance_sampler=UniformSplitSampler(p=1.0), + instance_sampler=transform.UniformSplitSampler(p=1.0), instance_length=train_length, prediction_length=pred_length, time_series_fields=["some_time_feature"], @@ -217,8 +200,6 @@ def test_CanonicalInstanceSplitter( use_prediction_features=use_prediction_features, ) - # assert_serializable(t) - other_feat = np.arange(len(target) + 100) data = { "start": start, @@ -256,39 +237,39 @@ def test_Transformation(): pred_length = 10 - t = Chain( + t = transform.Chain( trans=[ - AddTimeFeatures( + transform.AddTimeFeatures( start_field=FieldName.START, target_field=FieldName.TARGET, output_field="time_feat", time_features=[ - DayOfWeek(), - DayOfMonth(), - MonthOfYear(), + time_feature.DayOfWeek(), + time_feature.DayOfMonth(), + time_feature.MonthOfYear(), ], pred_length=pred_length, ), - AddAgeFeature( + transform.AddAgeFeature( target_field=FieldName.TARGET, output_field="age", pred_length=pred_length, log_scale=True, ), - AddObservedValuesIndicator( + transform.AddObservedValuesIndicator( target_field=FieldName.TARGET, output_field="observed_values" ), - VstackFeatures( + transform.VstackFeatures( output_field="dynamic_feat", input_fields=["age", "time_feat"], drop_inputs=True, ), - InstanceSplitter( + transform.InstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - train_sampler=ExpectedNumInstanceSampler(num_instances=4), + train_sampler=transform.ExpectedNumInstanceSampler(num_instances=4), past_length=train_length, future_length=pred_length, time_series_fields=["dynamic_feat", "observed_values"], @@ -296,8 +277,6 @@ def test_Transformation(): ] ) - # assert_serializable(t) - for u in t(iter(ds), is_train=True): print(u) @@ -323,51 +302,49 @@ def test_multi_dim_transformation(is_train): first_dim[-1] = np.nan second_dim[0] = np.nan - t = Chain( + t = transform.Chain( trans=[ - AddTimeFeatures( + transform.AddTimeFeatures( start_field=FieldName.START, target_field=FieldName.TARGET, output_field="time_feat", time_features=[ - DayOfWeek(), - DayOfMonth(), - MonthOfYear(), + time_feature.DayOfWeek(), + time_feature.DayOfMonth(), + time_feature.MonthOfYear(), ], pred_length=pred_length, ), - AddAgeFeature( + transform.AddAgeFeature( target_field=FieldName.TARGET, output_field="age", pred_length=pred_length, log_scale=True, ), - AddObservedValuesIndicator( + transform.AddObservedValuesIndicator( target_field=FieldName.TARGET, output_field="observed_values", convert_nans=False, ), - VstackFeatures( + transform.VstackFeatures( output_field="dynamic_feat", input_fields=["age", "time_feat"], drop_inputs=True, ), - InstanceSplitter( + transform.InstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - train_sampler=ExpectedNumInstanceSampler(num_instances=4), + train_sampler=transform.ExpectedNumInstanceSampler(num_instances=4), past_length=train_length, future_length=pred_length, time_series_fields=["dynamic_feat", "observed_values"], - batch_first=False, + output_NTC=False, ), ] ) - # assert_serializable(t) - if is_train: for u in t(iter(ds), is_train=True): assert_shape(u["past_target"], (2, 10)) @@ -406,14 +383,14 @@ def test_ExpectedNumInstanceSampler(): pred_length = 1 ds = make_dataset(N, train_length) - t = Chain( + t = transform.Chain( trans=[ - InstanceSplitter( + transform.InstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - train_sampler=ExpectedNumInstanceSampler(num_instances=4), + train_sampler=transform.ExpectedNumInstanceSampler(num_instances=4), past_length=train_length, future_length=pred_length, pick_incomplete=True, @@ -421,8 +398,6 @@ def test_ExpectedNumInstanceSampler(): ] ) - # assert_serializable(t) - scale_hist = ScaleHistogram() repetition = 2 @@ -446,14 +421,16 @@ def test_BucketInstanceSampler(): dataset_stats = calculate_dataset_statistics(ds) - t = Chain( + t = transform.Chain( trans=[ - InstanceSplitter( + transform.InstanceSplitter( target_field=FieldName.TARGET, is_pad_field=FieldName.IS_PAD, start_field=FieldName.START, forecast_start_field=FieldName.FORECAST_START, - train_sampler=BucketInstanceSampler(dataset_stats.scale_histogram), + train_sampler=transform.BucketInstanceSampler( + dataset_stats.scale_histogram + ), past_length=train_length, future_length=pred_length, pick_incomplete=True, @@ -461,8 +438,6 @@ def test_BucketInstanceSampler(): ] ) - # assert_serializable(t) - scale_hist = ScaleHistogram() repetition = 200 @@ -480,6 +455,292 @@ def test_BucketInstanceSampler(): assert abs(expected_values[i] - found_values[i] < expected_values[i] * 0.3) +def test_cdf_to_gaussian_transformation(): + def make_test_data(): + target = np.array( + [0, 0, 0, 0, 10, 10, 20, 20, 30, 30, 40, 50, 59, 60, 60, 70, 80, 90, 100,] + ).tolist() + + np.random.shuffle(target) + + multi_dim_target = np.array([target, target]).transpose() + + past_is_pad = np.array([[0] * len(target)]).transpose() + + past_observed_target = np.array( + [[1] * len(target), [1] * len(target)] + ).transpose() + + ds = ListDataset( + # Mimic output from InstanceSplitter + data_iter=[ + { + "start": "2012-01-01", + "target": multi_dim_target, + "past_target": multi_dim_target, + "future_target": multi_dim_target, + "past_is_pad": past_is_pad, + f"past_{FieldName.OBSERVED_VALUES}": past_observed_target, + } + ], + freq="1D", + one_dim_target=False, + ) + return ds + + def make_fake_output(u: DataEntry): + fake_output = np.expand_dims( + np.expand_dims(u["past_target_cdf"], axis=0), axis=0 + ) + return fake_output + + ds = make_test_data() + + t = transform.Chain( + trans=[ + transform.CDFtoGaussianTransform( + target_field=FieldName.TARGET, + observed_values_field=FieldName.OBSERVED_VALUES, + max_context_length=20, + target_dim=2, + ) + ] + ) + + for u in t(iter(ds), is_train=False): + + fake_output = make_fake_output(u) + + # Fake transformation chain output + u["past_target_sorted"] = torch.tensor( + np.expand_dims(u["past_target_sorted"], axis=0) + ) + + u["slopes"] = torch.tensor(np.expand_dims(u["slopes"], axis=0)) + + u["intercepts"] = torch.tensor(np.expand_dims(u["intercepts"], axis=0)) + + back_transformed = transform.cdf_to_gaussian_forward_transform(u, fake_output) + + # Get any sample/batch (slopes[i][:, d]they are all the same) + back_transformed = back_transformed[0][0] + + original_target = u["target"] + + # Original target and back-transformed target should be the same + assert np.allclose(original_target, back_transformed) + + +def test_gaussian_cdf(): + try: + from scipy.stats import norm + except: + pytest.skip("scipy not installed skipping test for erf") + + x = np.array( + [-1000, -100, -10] + np.linspace(-2, 2, 1001).tolist() + [10, 100, 1000] + ) + y_gluonts = transform.CDFtoGaussianTransform.standard_gaussian_cdf(x) + y_scipy = norm.cdf(x) + + assert np.allclose(y_gluonts, y_scipy, atol=1e-7) + + +def test_gaussian_ppf(): + try: + from scipy.stats import norm + except: + pytest.skip("scipy not installed skipping test for erf") + + x = np.linspace(0.0001, 0.9999, 1001) + y_gluonts = transform.CDFtoGaussianTransform.standard_gaussian_ppf(x) + y_scipy = norm.ppf(x) + + assert np.allclose(y_gluonts, y_scipy, atol=1e-7) + + +def test_target_dim_indicator(): + target = np.array([0, 2, 3, 10]).tolist() + + multi_dim_target = np.array([target, target, target, target]) + dataset = ListDataset( + data_iter=[{"start": "2012-01-01", "target": multi_dim_target}], + freq="1D", + one_dim_target=False, + ) + + t = transform.Chain( + trans=[ + transform.TargetDimIndicator( + target_field=FieldName.TARGET, field_name="target_dimensions" + ) + ] + ) + + for data_entry in t(dataset, is_train=True): + assert (data_entry["target_dimensions"] == np.array([0, 1, 2, 3])).all() + + +@pytest.fixture +def point_process_dataset(): + + ia_times = np.array([0.2, 0.7, 0.2, 0.5, 0.3, 0.3, 0.2, 0.1]) + marks = np.array([0, 1, 2, 0, 1, 2, 2, 2]) + + lds = ListDataset( + [ + { + "target": np.c_[ia_times, marks].T, + "start": pd.Timestamp("2011-01-01 00:00:00", freq="H"), + "end": pd.Timestamp("2011-01-01 03:00:00", freq="H"), + } + ], + freq="H", + one_dim_target=False, + ) + + return lds + + +class MockContinuousTimeSampler(transform.ContinuousTimePointSampler): + # noinspection PyMissingConstructor,PyUnusedLocal + def __init__(self, ret_values, *args, **kwargs): + self._ret_values = ret_values + + def __call__(self, *args, **kwargs): + return np.array(self._ret_values) + + +def test_ctsplitter_mask_sorted(point_process_dataset): + d = next(iter(point_process_dataset)) + + ia_times = d["target"][0, :] + + ts = np.cumsum(ia_times) + + splitter = transform.ContinuousTimeInstanceSplitter( + 2, 1, train_sampler=transform.ContinuousTimeUniformSampler(num_instances=10), + ) + + # no boundary conditions + res = splitter._mask_sorted(ts, 1, 2) + assert all([a == b for a, b in zip([2, 3, 4], res)]) + + # lower bound equal, exclusive of upper bound + res = splitter._mask_sorted(np.array([1, 2, 3, 4, 5, 6]), 1, 2) + assert all([a == b for a, b in zip([0], res)]) + + +def test_ctsplitter_no_train_last_point(point_process_dataset): + splitter = transform.ContinuousTimeInstanceSplitter( + 2, 1, train_sampler=transform.ContinuousTimeUniformSampler(num_instances=10), + ) + + iter_de = splitter(point_process_dataset, is_train=False) + + d_out = next(iter(iter_de)) + + assert "future_target" not in d_out + assert "future_valid_length" not in d_out + assert "past_target" in d_out + assert "past_valid_length" in d_out + + assert d_out["past_valid_length"] == 6 + assert np.allclose( + [0.1, 0.5, 0.3, 0.3, 0.2, 0.1], d_out["past_target"][..., 0], atol=0.01 + ) + + +def test_ctsplitter_train_correct(point_process_dataset): + splitter = transform.ContinuousTimeInstanceSplitter( + 1, + 1, + train_sampler=MockContinuousTimeSampler( + ret_values=[1.01, 1.5, 1.99], num_instances=3 + ), + ) + + iter_de = splitter(point_process_dataset, is_train=True) + + outputs = list(iter_de) + + assert outputs[0]["past_valid_length"] == 2 + assert outputs[0]["future_valid_length"] == 3 + + assert np.allclose(outputs[0]["past_target"], np.array([[0.19, 0.7], [0, 1]]).T) + assert np.allclose( + outputs[0]["future_target"], np.array([[0.09, 0.5, 0.3], [2, 0, 1]]).T + ) + + assert outputs[1]["past_valid_length"] == 2 + assert outputs[1]["future_valid_length"] == 4 + + assert outputs[2]["past_valid_length"] == 3 + assert outputs[2]["future_valid_length"] == 3 + + +def test_ctsplitter_train_correct_out_count(point_process_dataset): + + # produce new TPP data by shuffling existing TS instance + def shuffle_iterator(num_duplications=5): + for entry in point_process_dataset: + for i in range(num_duplications): + d = dict.copy(entry) + d["target"] = np.random.permutation(d["target"].T).T + yield d + + splitter = transform.ContinuousTimeInstanceSplitter( + 1, + 1, + train_sampler=MockContinuousTimeSampler( + ret_values=[1.01, 1.5, 1.99], num_instances=3 + ), + ) + + iter_de = splitter(shuffle_iterator(), is_train=True) + + outputs = list(iter_de) + + assert len(outputs) == 5 * 3 + + +def test_ctsplitter_train_samples_correct_times(point_process_dataset): + + splitter = transform.ContinuousTimeInstanceSplitter( + 1.25, 1.25, train_sampler=transform.ContinuousTimeUniformSampler(20) + ) + + iter_de = splitter(point_process_dataset, is_train=True) + + assert all( + [ + ( + pd.Timestamp("2011-01-01 01:15:00") + <= d["forecast_start"] + <= pd.Timestamp("2011-01-01 01:45:00") + ) + for d in iter_de + ] + ) + + +def test_ctsplitter_train_short_intervals(point_process_dataset): + splitter = transform.ContinuousTimeInstanceSplitter( + 0.01, + 0.01, + train_sampler=MockContinuousTimeSampler( + ret_values=[1.01, 1.5, 1.99], num_instances=3 + ), + ) + + iter_de = splitter(point_process_dataset, is_train=True) + + for d in iter_de: + assert d["future_valid_length"] == d["past_valid_length"] == 0 + assert np.prod(np.shape(d["past_target"])) == 0 + assert np.prod(np.shape(d["future_target"])) == 0 + + def make_dataset(N, train_length): # generates 2 ** N - 1 timeseries with constant increasing values n = 2 ** N - 1 @@ -500,6 +761,7 @@ def assert_shape(array: np.array, reference_shape: Tuple[int, int]): array.shape == reference_shape ), f"Shape should be {reference_shape} but found {array.shape}." + def assert_padded_array( sampled_array: np.array, reference_array: np.array, padding_array: np.array ):