From 73f51769f8fd808c036b58684364c8e611d010ff Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Mon, 10 Jun 2019 13:33:31 -0700 Subject: [PATCH 01/12] Initial creation of the IidSpikeDetector files to see what works and what doesn't. --- src/python/nimbusml.pyproj | 7 ++ .../IidSpikeDetector_df.py | 12 +++ .../internal/core/time_series/__init__.py | 0 .../core/time_series/iidspikedetector.py | 91 +++++++++++++++++++ src/python/nimbusml/time_series/__init__.py | 5 + .../nimbusml/time_series/iidspikedetector.py | 91 +++++++++++++++++++ src/python/tools/manifest_diff.json | 6 ++ 7 files changed, 212 insertions(+) create mode 100644 src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py create mode 100644 src/python/nimbusml/internal/core/time_series/__init__.py create mode 100644 src/python/nimbusml/internal/core/time_series/iidspikedetector.py create mode 100644 src/python/nimbusml/time_series/__init__.py create mode 100644 src/python/nimbusml/time_series/iidspikedetector.py diff --git a/src/python/nimbusml.pyproj b/src/python/nimbusml.pyproj index 2fecda2d..5e807454 100644 --- a/src/python/nimbusml.pyproj +++ b/src/python/nimbusml.pyproj @@ -109,6 +109,7 @@ + @@ -224,6 +225,8 @@ + + @@ -571,6 +574,8 @@ + + @@ -743,6 +748,7 @@ + @@ -780,6 +786,7 @@ + diff --git a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py new file mode 100644 index 00000000..7bef3e2e --- /dev/null +++ b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py @@ -0,0 +1,12 @@ +############################################################################### +# IidSpikeDetector +import numpy as np +import pandas as pd +from nimbusml.time_series import IidSpikeDetector + +X_train = pd.Series([5,5,5,5,5,10,5,5,5,5,5]) + +isd = IidSpikeDetector() +result = isd.fit_transform(X_train) + +print(result) diff --git a/src/python/nimbusml/internal/core/time_series/__init__.py b/src/python/nimbusml/internal/core/time_series/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/src/python/nimbusml/internal/core/time_series/iidspikedetector.py b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py new file mode 100644 index 00000000..00712d77 --- /dev/null +++ b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py @@ -0,0 +1,91 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +IidSpikeDetector +""" + +__all__ = ["IidSpikeDetector"] + + +from ...entrypoints.timeseriesprocessingentrypoints_iidspikedetector import \ + timeseriesprocessingentrypoints_iidspikedetector +from ...utils.utils import trace +from ..base_pipeline_item import BasePipelineItem, DefaultSignature + + +class IidSpikeDetector(BasePipelineItem, DefaultSignature): + """ + + This transform detects the spikes in a i.i.d. sequence using adaptive + kernel density estimation. + + .. remarks:: + ``IIDSpikeDetector`` assumes a sequence of data points that are + independently sampled from one stationary + distribution. `Adaptive kernel density estimation + `_ + is used to model the distribution. + The `p-value score + indicates the likelihood of the current observation according to + the estimated distribution. The lower its value, the more likely the + current point is an outlier. + + :param confidence: The confidence for spike detection in the range [0, + 100]. + + :param side: The argument that determines whether to detect positive or + negative anomalies, or both. Available options are {``Positive``, + ``Negative``, ``TwoSided``}. + + :param pvalue_history_length: The size of the sliding window for computing + the p-value. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`SsaSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/IidSpikePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + confidence=99.0, + side='TwoSided', + pvalue_history_length=100, + **params): + BasePipelineItem.__init__( + self, type='transform', **params) + + self.confidence = confidence + self.side = side + self.pvalue_history_length = pvalue_history_length + + @property + def _entrypoint(self): + return timeseriesprocessingentrypoints_iidspikedetector + + @trace + def _get_node(self, **all_args): + algo_args = dict( + source=self.source, + name=self._name_or_source, + confidence=self.confidence, + side=self.side, + pvalue_history_length=self.pvalue_history_length) + + all_args.update(algo_args) + return self._entrypoint(**all_args) diff --git a/src/python/nimbusml/time_series/__init__.py b/src/python/nimbusml/time_series/__init__.py new file mode 100644 index 00000000..60f1cd18 --- /dev/null +++ b/src/python/nimbusml/time_series/__init__.py @@ -0,0 +1,5 @@ +from .iidspikedetector import IidSpikeDetector + +__all__ = [ + 'IidSpikeDetector' +] diff --git a/src/python/nimbusml/time_series/iidspikedetector.py b/src/python/nimbusml/time_series/iidspikedetector.py new file mode 100644 index 00000000..acb6dbaa --- /dev/null +++ b/src/python/nimbusml/time_series/iidspikedetector.py @@ -0,0 +1,91 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +IidSpikeDetector +""" + +__all__ = ["IidSpikeDetector"] + + +from sklearn.base import TransformerMixin + +from ..base_transform import BaseTransform +from ..internal.core.time_series.iidspikedetector import \ + IidSpikeDetector as core +from ..internal.utils.utils import trace + + +class IidSpikeDetector(core, BaseTransform, TransformerMixin): + """ + + This transform detects the spikes in a i.i.d. sequence using adaptive + kernel density estimation. + + .. remarks:: + ``IIDSpikeDetector`` assumes a sequence of data points that are + independently sampled from one stationary + distribution. `Adaptive kernel density estimation + `_ + is used to model the distribution. + The `p-value score + indicates the likelihood of the current observation according to + the estimated distribution. The lower its value, the more likely the + current point is an outlier. + + :param columns: see `Columns `_. + + :param confidence: The confidence for spike detection in the range [0, + 100]. + + :param side: The argument that determines whether to detect positive or + negative anomalies, or both. Available options are {``Positive``, + ``Negative``, ``TwoSided``}. + + :param pvalue_history_length: The size of the sliding window for computing + the p-value. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`SsaSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/IidSpikePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + confidence=99.0, + side='TwoSided', + pvalue_history_length=100, + columns=None, + **params): + + if columns: + params['columns'] = columns + BaseTransform.__init__(self, **params) + core.__init__( + self, + confidence=confidence, + side=side, + pvalue_history_length=pvalue_history_length, + **params) + self._columns = columns + + def get_params(self, deep=False): + """ + Get the parameters for this operator. + """ + return core.get_params(self) diff --git a/src/python/tools/manifest_diff.json b/src/python/tools/manifest_diff.json index c19aad98..330a49f2 100644 --- a/src/python/tools/manifest_diff.json +++ b/src/python/tools/manifest_diff.json @@ -539,6 +539,12 @@ "Module": "decomposition", "Type": "Anomaly" }, + { + "Name": "TimeSeriesProcessingEntryPoints.IidSpikeDetector", + "NewName": "IidSpikeDetector", + "Module": "time_series", + "Type": "Transform" + }, { "Name": "Trainers.PoissonRegressor", "NewName": "PoissonRegressionRegressor", From 9ac77c7da57ec765a13646d0db36531a6917b669 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Mon, 10 Jun 2019 14:04:37 -0700 Subject: [PATCH 02/12] Import the Microsoft.ML.TimeSeries assembly in to the project. --- src/DotNetBridge/Bridge.cs | 2 ++ src/DotNetBridge/DotNetBridge.csproj | 1 + src/Platforms/build.csproj | 1 + 3 files changed, 4 insertions(+) diff --git a/src/DotNetBridge/Bridge.cs b/src/DotNetBridge/Bridge.cs index 1395c998..26e5a84d 100644 --- a/src/DotNetBridge/Bridge.cs +++ b/src/DotNetBridge/Bridge.cs @@ -17,6 +17,7 @@ using Microsoft.ML.Trainers.FastTree; using Microsoft.ML.Trainers.LightGbm; using Microsoft.ML.Transforms; +using Microsoft.ML.TimeSeries; namespace Microsoft.MachineLearning.DotNetBridge { @@ -328,6 +329,7 @@ private static unsafe int GenericExec(EnvironmentBlock* penv, sbyte* psz, int cd //env.ComponentCatalog.RegisterAssembly(typeof(SaveOnnxCommand).Assembly); //env.ComponentCatalog.RegisterAssembly(typeof(TimeSeriesProcessingEntryPoints).Assembly); //env.ComponentCatalog.RegisterAssembly(typeof(ParquetLoader).Assembly); + env.ComponentCatalog.RegisterAssembly(typeof(ForecastExtensions).Assembly); using (var ch = host.Start("Executing")) { diff --git a/src/DotNetBridge/DotNetBridge.csproj b/src/DotNetBridge/DotNetBridge.csproj index 92365878..fab49e2e 100644 --- a/src/DotNetBridge/DotNetBridge.csproj +++ b/src/DotNetBridge/DotNetBridge.csproj @@ -41,5 +41,6 @@ + diff --git a/src/Platforms/build.csproj b/src/Platforms/build.csproj index 7491fac8..e75aa8f3 100644 --- a/src/Platforms/build.csproj +++ b/src/Platforms/build.csproj @@ -20,6 +20,7 @@ + From 0e55ae3679d442ae238329cce9b4a65bd08dd88f Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Tue, 11 Jun 2019 08:42:20 -0700 Subject: [PATCH 03/12] Use 'PassAs' in manifest.json to fix the source parameter name. --- .../core/time_series/iidspikedetector.py | 19 ++++++++++++++++++- ...sprocessingentrypoints_iidspikedetector.py | 8 ++++---- src/python/tools/manifest.json | 3 ++- 3 files changed, 24 insertions(+), 6 deletions(-) diff --git a/src/python/nimbusml/internal/core/time_series/iidspikedetector.py b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py index 00712d77..31c7ad4f 100644 --- a/src/python/nimbusml/internal/core/time_series/iidspikedetector.py +++ b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py @@ -80,8 +80,25 @@ def _entrypoint(self): @trace def _get_node(self, **all_args): + + input_column = self.input + if input_column is None and 'input' in all_args: + input_column = all_args['input'] + if 'input' in all_args: + all_args.pop('input') + + # validate input + if input_column is None: + raise ValueError( + "'None' input passed when it cannot be none.") + + if not isinstance(input_column, str): + raise ValueError( + "input has to be a string, instead got %s" % + type(input_column)) + algo_args = dict( - source=self.source, + column=input_column, name=self._name_or_source, confidence=self.confidence, side=self.side, diff --git a/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py b/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py index 113ddc72..0f2cf5ed 100644 --- a/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py +++ b/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py @@ -10,7 +10,7 @@ def timeseriesprocessingentrypoints_iidspikedetector( - source, + column, data, name, output_data=None, @@ -24,7 +24,7 @@ def timeseriesprocessingentrypoints_iidspikedetector( This transform detects the spikes in a i.i.d. sequence using adaptive kernel density estimation. - :param source: The name of the source column. (inputs). + :param column: The name of the source column. (inputs). :param data: Input dataset (inputs). :param name: The name of the new column. (inputs). :param confidence: The confidence for spike detection in the @@ -41,9 +41,9 @@ def timeseriesprocessingentrypoints_iidspikedetector( inputs = {} outputs = {} - if source is not None: + if column is not None: inputs['Source'] = try_set( - obj=source, + obj=column, none_acceptable=False, is_of_type=str, is_column=True) diff --git a/src/python/tools/manifest.json b/src/python/tools/manifest.json index 67951c74..e01daf77 100644 --- a/src/python/tools/manifest.json +++ b/src/python/tools/manifest.json @@ -3489,7 +3489,8 @@ "ShortName": "ispike", "Inputs": [ { - "Name": "Source", + "Name": "Column", + "PassAs": "Source", "Type": "String", "Desc": "The name of the source column.", "Aliases": [ From 204128bf9de54f55c6328881cc490ed29df48d61 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Tue, 11 Jun 2019 08:45:32 -0700 Subject: [PATCH 04/12] Use float32 for data dtype in IidSpikeDetector example. --- .../examples/examples_from_dataframe/IidSpikeDetector_df.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py index 7bef3e2e..bb8105e7 100644 --- a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py +++ b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py @@ -4,9 +4,9 @@ import pandas as pd from nimbusml.time_series import IidSpikeDetector -X_train = pd.Series([5,5,5,5,5,10,5,5,5,5,5]) +X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], dtype=np.float32) isd = IidSpikeDetector() -result = isd.fit_transform(X_train) +result = isd.fit_transform(X_train, verbose=1) print(result) From 0a60da151150364006e0134f3cfeac2a95c3e140 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Tue, 11 Jun 2019 15:38:06 -0700 Subject: [PATCH 05/12] Convert IidSpikeDetector to a standard transform. Add examples and tests. --- src/python/nimbusml.pyproj | 4 ++ .../nimbusml/examples/IidSpikeDetector.py | 39 ++++++++++++++++ .../IidSpikeDetector_df.py | 23 ++++++++-- .../core/time_series/iidspikedetector.py | 19 +------- ...sprocessingentrypoints_iidspikedetector.py | 8 ++-- .../nimbusml/tests/time_series/__init__.py | 0 .../time_series/test_iidspikedetector.py | 46 +++++++++++++++++++ src/python/tools/manifest.json | 3 +- 8 files changed, 114 insertions(+), 28 deletions(-) create mode 100644 src/python/nimbusml/examples/IidSpikeDetector.py create mode 100644 src/python/nimbusml/tests/time_series/__init__.py create mode 100644 src/python/nimbusml/tests/time_series/test_iidspikedetector.py diff --git a/src/python/nimbusml.pyproj b/src/python/nimbusml.pyproj index 5e807454..64764777 100644 --- a/src/python/nimbusml.pyproj +++ b/src/python/nimbusml.pyproj @@ -160,6 +160,7 @@ + @@ -574,6 +575,8 @@ + + @@ -770,6 +773,7 @@ + diff --git a/src/python/nimbusml/examples/IidSpikeDetector.py b/src/python/nimbusml/examples/IidSpikeDetector.py new file mode 100644 index 00000000..790c320d --- /dev/null +++ b/src/python/nimbusml/examples/IidSpikeDetector.py @@ -0,0 +1,39 @@ +############################################################################### +# IidSpikeDetector +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import IidSpikeDetector +from nimbusml.preprocessing.schema import TypeConverter + +# data input (as a FileDataStream) +path = get_dataset('timeseries').as_filepath() + +data = FileDataStream.read_csv(path) +print(data.head()) +# t1 t2 t3 +# 0 0.01 0.01 0.0100 +# 1 0.02 0.02 0.0200 +# 2 0.03 0.03 0.0200 +# 3 0.03 0.03 0.0250 +# 4 0.03 0.03 0.0005 + +# define the training pipeline +pipeline = Pipeline([ + TypeConverter(result_type='R4'), + IidSpikeDetector(columns={'t2_spikes': 't2'}, pvalue_history_length=5) +]) + +result = pipeline.fit_transform(data) +print(result) +# t1 t2 t3 t2_spikes.Alert t2_spikes.Raw Score t2_spikes.P-Value Score +# 0 0.01 0.01 0.0100 0.0 0.01 5.000000e-01 +# 1 0.02 0.02 0.0200 0.0 0.02 4.960106e-01 +# 2 0.03 0.03 0.0200 0.0 0.03 1.139087e-02 +# 3 0.03 0.03 0.0250 0.0 0.03 2.058296e-01 +# 4 0.03 0.03 0.0005 0.0 0.03 2.804577e-01 +# 5 0.03 0.05 0.0100 1.0 0.05 3.743552e-03 +# 6 0.05 0.07 0.0500 1.0 0.07 4.136079e-03 +# 7 0.07 0.09 0.0900 0.0 0.09 2.242496e-02 +# 8 0.09 99.00 99.0000 1.0 99.00 1.000000e-08 +# 9 1.10 0.10 0.1000 0.0 0.10 4.015681e-01 + diff --git a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py index bb8105e7..88b50565 100644 --- a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py +++ b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py @@ -4,9 +4,24 @@ import pandas as pd from nimbusml.time_series import IidSpikeDetector -X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], dtype=np.float32) +X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts", dtype=np.float32) -isd = IidSpikeDetector() -result = isd.fit_transform(X_train, verbose=1) +isd = IidSpikeDetector(confidence=95, pvalue_history_length=2.5) << {'result': 'ts'} -print(result) +isd.fit(X_train, verbose=1) +data = isd.transform(X_train) + +print(data) + +# ts result.Alert result.Raw Score result.P-Value Score +# 0 5.0 0.0 5.0 5.000000e-01 +# 1 5.0 0.0 5.0 5.000000e-01 +# 2 5.0 0.0 5.0 5.000000e-01 +# 3 5.0 0.0 5.0 5.000000e-01 +# 4 5.0 0.0 5.0 5.000000e-01 +# 5 10.0 1.0 10.0 1.000000e-08 +# 6 5.0 0.0 5.0 2.613750e-01 +# 7 5.0 0.0 5.0 2.613750e-01 +# 8 5.0 0.0 5.0 5.000000e-01 +# 9 5.0 0.0 5.0 5.000000e-01 +# 10 5.0 0.0 5.0 5.000000e-01 diff --git a/src/python/nimbusml/internal/core/time_series/iidspikedetector.py b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py index 31c7ad4f..00712d77 100644 --- a/src/python/nimbusml/internal/core/time_series/iidspikedetector.py +++ b/src/python/nimbusml/internal/core/time_series/iidspikedetector.py @@ -80,25 +80,8 @@ def _entrypoint(self): @trace def _get_node(self, **all_args): - - input_column = self.input - if input_column is None and 'input' in all_args: - input_column = all_args['input'] - if 'input' in all_args: - all_args.pop('input') - - # validate input - if input_column is None: - raise ValueError( - "'None' input passed when it cannot be none.") - - if not isinstance(input_column, str): - raise ValueError( - "input has to be a string, instead got %s" % - type(input_column)) - algo_args = dict( - column=input_column, + source=self.source, name=self._name_or_source, confidence=self.confidence, side=self.side, diff --git a/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py b/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py index 0f2cf5ed..113ddc72 100644 --- a/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py +++ b/src/python/nimbusml/internal/entrypoints/timeseriesprocessingentrypoints_iidspikedetector.py @@ -10,7 +10,7 @@ def timeseriesprocessingentrypoints_iidspikedetector( - column, + source, data, name, output_data=None, @@ -24,7 +24,7 @@ def timeseriesprocessingentrypoints_iidspikedetector( This transform detects the spikes in a i.i.d. sequence using adaptive kernel density estimation. - :param column: The name of the source column. (inputs). + :param source: The name of the source column. (inputs). :param data: Input dataset (inputs). :param name: The name of the new column. (inputs). :param confidence: The confidence for spike detection in the @@ -41,9 +41,9 @@ def timeseriesprocessingentrypoints_iidspikedetector( inputs = {} outputs = {} - if column is not None: + if source is not None: inputs['Source'] = try_set( - obj=column, + obj=source, none_acceptable=False, is_of_type=str, is_column=True) diff --git a/src/python/nimbusml/tests/time_series/__init__.py b/src/python/nimbusml/tests/time_series/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/src/python/nimbusml/tests/time_series/test_iidspikedetector.py b/src/python/nimbusml/tests/time_series/test_iidspikedetector.py new file mode 100644 index 00000000..9c37e456 --- /dev/null +++ b/src/python/nimbusml/tests/time_series/test_iidspikedetector.py @@ -0,0 +1,46 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- + +import unittest + +import numpy as np +import pandas as pd +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import IidSpikeDetector +from nimbusml.preprocessing.schema import TypeConverter + + +class TestIidSpikeDetector(unittest.TestCase): + + def test_correct_data_is_marked_as_anomaly(self): + X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts", dtype=np.float32) + isd = IidSpikeDetector(confidence=95, pvalue_history_length=3) << {'result': 'ts'} + data = isd.fit_transform(X_train) + + data = data.loc[data['result.Alert'] == 1.0] + self.assertEqual(len(data), 1) + self.assertEqual(data.iloc[0]['ts'], 10.0) + + def test_multiple_user_specified_columns_is_not_allowed(self): + path = get_dataset('timeseries').as_filepath() + data = FileDataStream.read_csv(path) + + try: + pipeline = Pipeline([ + TypeConverter(result_type='R4'), + IidSpikeDetector(columns=['t2', 't3'], pvalue_history_length=5) + ]) + pipeline.fit_transform(data) + + except RuntimeError as e: + self.assertTrue('Only one column is allowed' in str(e)) + return + + self.fail() + + +if __name__ == '__main__': + unittest.main() diff --git a/src/python/tools/manifest.json b/src/python/tools/manifest.json index e01daf77..67951c74 100644 --- a/src/python/tools/manifest.json +++ b/src/python/tools/manifest.json @@ -3489,8 +3489,7 @@ "ShortName": "ispike", "Inputs": [ { - "Name": "Column", - "PassAs": "Source", + "Name": "Source", "Type": "String", "Desc": "The name of the source column.", "Aliases": [ From 33043a9d7ad1080d726c51c1644e2c05737a0940 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Wed, 12 Jun 2019 10:28:51 -0700 Subject: [PATCH 06/12] Add pre-transform to IidSpikeDetector to fix incompatible data types. --- .../nimbusml/examples/IidSpikeDetector.py | 1 - .../IidSpikeDetector_df.py | 2 +- .../time_series/test_iidspikedetector.py | 21 +++++++++++++++++-- .../nimbusml/time_series/iidspikedetector.py | 10 +++++++++ src/python/tools/compiler_utils.py | 6 ++++++ 5 files changed, 36 insertions(+), 4 deletions(-) diff --git a/src/python/nimbusml/examples/IidSpikeDetector.py b/src/python/nimbusml/examples/IidSpikeDetector.py index 790c320d..1496c66e 100644 --- a/src/python/nimbusml/examples/IidSpikeDetector.py +++ b/src/python/nimbusml/examples/IidSpikeDetector.py @@ -19,7 +19,6 @@ # define the training pipeline pipeline = Pipeline([ - TypeConverter(result_type='R4'), IidSpikeDetector(columns={'t2_spikes': 't2'}, pvalue_history_length=5) ]) diff --git a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py index 88b50565..723f7b2a 100644 --- a/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py +++ b/src/python/nimbusml/examples/examples_from_dataframe/IidSpikeDetector_df.py @@ -4,7 +4,7 @@ import pandas as pd from nimbusml.time_series import IidSpikeDetector -X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts", dtype=np.float32) +X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts") isd = IidSpikeDetector(confidence=95, pvalue_history_length=2.5) << {'result': 'ts'} diff --git a/src/python/nimbusml/tests/time_series/test_iidspikedetector.py b/src/python/nimbusml/tests/time_series/test_iidspikedetector.py index 9c37e456..6ef5ac89 100644 --- a/src/python/nimbusml/tests/time_series/test_iidspikedetector.py +++ b/src/python/nimbusml/tests/time_series/test_iidspikedetector.py @@ -16,7 +16,7 @@ class TestIidSpikeDetector(unittest.TestCase): def test_correct_data_is_marked_as_anomaly(self): - X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts", dtype=np.float32) + X_train = pd.Series([5, 5, 5, 5, 5, 10, 5, 5, 5, 5, 5], name="ts") isd = IidSpikeDetector(confidence=95, pvalue_history_length=3) << {'result': 'ts'} data = isd.fit_transform(X_train) @@ -30,7 +30,6 @@ def test_multiple_user_specified_columns_is_not_allowed(self): try: pipeline = Pipeline([ - TypeConverter(result_type='R4'), IidSpikeDetector(columns=['t2', 't3'], pvalue_history_length=5) ]) pipeline.fit_transform(data) @@ -41,6 +40,24 @@ def test_multiple_user_specified_columns_is_not_allowed(self): self.fail() + def test_pre_transform_does_not_convert_non_time_series_columns(self): + X_train = pd.DataFrame({ + 'Date': ['2017-01', '2017-02', '2017-03'], + 'Values': [5.0, 5.0, 5.0]}) + + self.assertEqual(len(X_train.dtypes), 2) + self.assertEqual(str(X_train.dtypes[0]), 'object') + self.assertTrue(str(X_train.dtypes[1]).startswith('float')) + + isd = IidSpikeDetector(confidence=95, pvalue_history_length=3) << 'Values' + data = isd.fit_transform(X_train) + + self.assertEqual(len(data.dtypes), 4) + self.assertEqual(str(data.dtypes[0]), 'object') + self.assertTrue(str(data.dtypes[1]).startswith('float')) + self.assertTrue(str(data.dtypes[2]).startswith('float')) + self.assertTrue(str(data.dtypes[3]).startswith('float')) + if __name__ == '__main__': unittest.main() diff --git a/src/python/nimbusml/time_series/iidspikedetector.py b/src/python/nimbusml/time_series/iidspikedetector.py index acb6dbaa..7f570003 100644 --- a/src/python/nimbusml/time_series/iidspikedetector.py +++ b/src/python/nimbusml/time_series/iidspikedetector.py @@ -89,3 +89,13 @@ def get_params(self, deep=False): Get the parameters for this operator. """ return core.get_params(self) + + def _nodes_with_presteps(self): + """ + Inserts preprocessing before this one. + """ + from ..preprocessing.schema import TypeConverter + return [ + TypeConverter( + result_type='R4')._steal_io(self), + self] diff --git a/src/python/tools/compiler_utils.py b/src/python/tools/compiler_utils.py index d7462c78..25faa82b 100644 --- a/src/python/tools/compiler_utils.py +++ b/src/python/tools/compiler_utils.py @@ -120,6 +120,10 @@ def _nodes_with_presteps(self): '''from ..schema import TypeConverter return [TypeConverter(result_type='R4')._steal_io(self), self]''' +timeseries_to_r4_converter = \ + '''from ..preprocessing.schema import TypeConverter +return [TypeConverter(result_type='R4')._steal_io(self), self]''' + _presteps = { 'MinMaxScaler': int_to_r4_converter, 'MeanVarianceScaler': int_to_r4_converter, @@ -127,6 +131,8 @@ def _nodes_with_presteps(self): 'Binner': int_to_r4_converter, # 'SupervisedBinner': int_to_r4_converter, # not exist in nimbusml + 'IidSpikeDetector': timeseries_to_r4_converter, + 'PcaTransformer': '''from ..preprocessing.schema import TypeConverter if type(self._columns) == dict: From e92ffe1a9b30b8e0b8564a2d41106435023260a0 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Wed, 12 Jun 2019 16:00:48 -0700 Subject: [PATCH 07/12] Fix issues with the test_estimator_checks IidSpikeDetector tests. --- src/python/tests/test_estimator_checks.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/python/tests/test_estimator_checks.py b/src/python/tests/test_estimator_checks.py index 07b1453c..1d36dd4d 100644 --- a/src/python/tests/test_estimator_checks.py +++ b/src/python/tests/test_estimator_checks.py @@ -16,6 +16,7 @@ from nimbusml.internal.entrypoints._ngramextractor_ngram import n_gram from nimbusml.preprocessing import TensorFlowScorer from nimbusml.preprocessing.filter import SkipFilter, TakeFilter +from nimbusml.time_series import IidSpikeDetector from sklearn.utils.estimator_checks import _yield_all_checks, MULTI_OUTPUT this = os.path.abspath(os.path.dirname(__file__)) @@ -53,6 +54,8 @@ # fix pending in PR, bug cant handle csr matrix 'RangeFilter': 'check_estimators_dtypes, ' 'check_estimator_sparse_data', + # time series does not currently support sparse matrices + 'IidSpikeDetector': 'check_estimator_sparse_data', # bug, low tolerance 'FastLinearRegressor': 'check_supervised_y_2d, ' 'check_regressor_data_not_an_array, ' @@ -180,6 +183,7 @@ 'NGramFeaturizer': NGramFeaturizer(word_feature_extractor=n_gram()), 'SkipFilter': SkipFilter(count=5), 'TakeFilter': TakeFilter(count=100000), + 'IidSpikeDetector': IidSpikeDetector(columns=['F0']), 'TensorFlowScorer': TensorFlowScorer( model_location=os.path.join( this, From 5b92c2ac42ec6ece5c8aeeb579d0a7e29bcadb45 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Fri, 14 Jun 2019 08:51:39 -0700 Subject: [PATCH 08/12] Remove unnecessary TypeConverter import in IidSpikeDetector example. --- src/python/nimbusml/examples/IidSpikeDetector.py | 1 - 1 file changed, 1 deletion(-) diff --git a/src/python/nimbusml/examples/IidSpikeDetector.py b/src/python/nimbusml/examples/IidSpikeDetector.py index 1496c66e..6876375f 100644 --- a/src/python/nimbusml/examples/IidSpikeDetector.py +++ b/src/python/nimbusml/examples/IidSpikeDetector.py @@ -3,7 +3,6 @@ from nimbusml import Pipeline, FileDataStream from nimbusml.datasets import get_dataset from nimbusml.time_series import IidSpikeDetector -from nimbusml.preprocessing.schema import TypeConverter # data input (as a FileDataStream) path = get_dataset('timeseries').as_filepath() From 9d9aa605e466f339e25d793095e076496c34a850 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Fri, 14 Jun 2019 10:53:17 -0700 Subject: [PATCH 09/12] Initial implementation of IidChangePointDetector. --- src/python/nimbusml.pyproj | 5 + .../examples/IidChangePointDetector.py | 38 ++++++ .../IidChangePointDetector_df.py | 34 +++++ .../time_series/iidchangepointdetector.py | 107 ++++++++++++++++ .../test_iidchangepointdetector.py | 48 +++++++ src/python/nimbusml/time_series/__init__.py | 4 +- .../time_series/iidchangepointdetector.py | 119 ++++++++++++++++++ src/python/tests/test_estimator_checks.py | 6 +- src/python/tools/compiler_utils.py | 1 + src/python/tools/manifest_diff.json | 6 + 10 files changed, 365 insertions(+), 3 deletions(-) create mode 100644 src/python/nimbusml/examples/IidChangePointDetector.py create mode 100644 src/python/nimbusml/examples/examples_from_dataframe/IidChangePointDetector_df.py create mode 100644 src/python/nimbusml/internal/core/time_series/iidchangepointdetector.py create mode 100644 src/python/nimbusml/tests/time_series/test_iidchangepointdetector.py create mode 100644 src/python/nimbusml/time_series/iidchangepointdetector.py diff --git a/src/python/nimbusml.pyproj b/src/python/nimbusml.pyproj index 64764777..e1924f93 100644 --- a/src/python/nimbusml.pyproj +++ b/src/python/nimbusml.pyproj @@ -88,6 +88,7 @@ + @@ -135,6 +136,7 @@ + @@ -226,6 +228,7 @@ + @@ -575,8 +578,10 @@ + + diff --git a/src/python/nimbusml/examples/IidChangePointDetector.py b/src/python/nimbusml/examples/IidChangePointDetector.py new file mode 100644 index 00000000..d8f9f4d8 --- /dev/null +++ b/src/python/nimbusml/examples/IidChangePointDetector.py @@ -0,0 +1,38 @@ +############################################################################### +# IidChangePointDetector +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import IidChangePointDetector + +# data input (as a FileDataStream) +path = get_dataset('timeseries').as_filepath() + +data = FileDataStream.read_csv(path) +print(data.head()) +# t1 t2 t3 +# 0 0.01 0.01 0.0100 +# 1 0.02 0.02 0.0200 +# 2 0.03 0.03 0.0200 +# 3 0.03 0.03 0.0250 +# 4 0.03 0.03 0.0005 + +# define the training pipeline +pipeline = Pipeline([ + IidChangePointDetector(columns={'t2_cp': 't2'}, change_history_length=4) +]) + +result = pipeline.fit_transform(data) +print(result) + +# t1 t2 t3 t2_cp.Alert t2_cp.Raw Score t2_cp.P-Value Score t2_cp.Martingale Score +# 0 0.01 0.01 0.0100 0.0 0.01 5.000000e-01 1.212573e-03 +# 1 0.02 0.02 0.0200 0.0 0.02 4.960106e-01 1.221347e-03 +# 2 0.03 0.03 0.0200 0.0 0.03 1.139087e-02 3.672914e-02 +# 3 0.03 0.03 0.0250 0.0 0.03 2.058296e-01 8.164447e-02 +# 4 0.03 0.03 0.0005 0.0 0.03 2.804577e-01 1.373786e-01 +# 5 0.03 0.05 0.0100 1.0 0.05 1.448886e-06 1.315014e+04 +# 6 0.05 0.07 0.0500 0.0 0.07 2.616611e-03 4.941587e+04 +# 7 0.07 0.09 0.0900 0.0 0.09 3.053187e-02 2.752614e+05 +# 8 0.09 99.00 99.0000 0.0 99.00 1.000000e-08 1.389396e+12 +# 9 1.10 0.10 0.1000 1.0 0.10 3.778296e-01 1.854344e+07 + diff --git a/src/python/nimbusml/examples/examples_from_dataframe/IidChangePointDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/IidChangePointDetector_df.py new file mode 100644 index 00000000..2401f118 --- /dev/null +++ b/src/python/nimbusml/examples/examples_from_dataframe/IidChangePointDetector_df.py @@ -0,0 +1,34 @@ +############################################################################### +# IidChangePointDetector +import pandas as pd +from nimbusml.time_series import IidChangePointDetector + +# Create a sample series with a change +input_data = [5, 5, 5, 5, 5, 5, 5, 5] +input_data.extend([7, 7, 7, 7, 7, 7, 7, 7]) + +X_train = pd.Series(input_data, name="ts") + +cpd = IidChangePointDetector(confidence=95, change_history_length=4) << {'result': 'ts'} +data = cpd.fit_transform(X_train) + +print(data) + +# ts result.Alert result.Raw Score result.P-Value Score result.Martingale Score +# 0 5 0.0 5.0 5.000000e-01 0.001213 +# 1 5 0.0 5.0 5.000000e-01 0.001213 +# 2 5 0.0 5.0 5.000000e-01 0.001213 +# 3 5 0.0 5.0 5.000000e-01 0.001213 +# 4 5 0.0 5.0 5.000000e-01 0.001213 +# 5 5 0.0 5.0 5.000000e-01 0.001213 +# 6 5 0.0 5.0 5.000000e-01 0.001213 +# 7 5 0.0 5.0 5.000000e-01 0.001213 +# 8 7 1.0 7.0 1.000000e-08 10298.666376 <-- alert is on, predicted changepoint +# 9 7 0.0 7.0 1.328455e-01 33950.164799 +# 10 7 0.0 7.0 2.613750e-01 60866.342063 +# 11 7 0.0 7.0 3.776152e-01 78362.038772 +# 12 7 0.0 7.0 5.000000e-01 0.009226 +# 13 7 0.0 7.0 5.000000e-01 0.002799 +# 14 7 0.0 7.0 5.000000e-01 0.001561 +# 15 7 0.0 7.0 5.000000e-01 0.001213 + diff --git a/src/python/nimbusml/internal/core/time_series/iidchangepointdetector.py b/src/python/nimbusml/internal/core/time_series/iidchangepointdetector.py new file mode 100644 index 00000000..ae874a1c --- /dev/null +++ b/src/python/nimbusml/internal/core/time_series/iidchangepointdetector.py @@ -0,0 +1,107 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +IidChangePointDetector +""" + +__all__ = ["IidChangePointDetector"] + + +from ...entrypoints.timeseriesprocessingentrypoints_iidchangepointdetector import \ + timeseriesprocessingentrypoints_iidchangepointdetector +from ...utils.utils import trace +from ..base_pipeline_item import BasePipelineItem, DefaultSignature + + +class IidChangePointDetector(BasePipelineItem, DefaultSignature): + """ + + This transform detects the change-points in an i.i.d. sequence using + adaptive kernel density estimation and martingales. + + .. remarks:: + ``IIDChangePointDetector`` assumes a sequence of data points that are + independently sampled from one + stationary distribution. `Adaptive kernel density estimation + `_ + is used to model the distribution. + + This transform detects + change points by calculating the martingale score for the sliding + window based on the estimated distribution. + The idea is based on the `Exchangeability + Martingales `_ that + detects a change of distribution over a stream of i.i.d. values. In + short, the value of the + martingale score starts increasing significantly when a sequence of + small p-values are detected in a row; this + indicates the change of the distribution of the underlying data + generation process. + + :param confidence: The confidence for change point detection in the range + [0, 100]. Used to set the threshold of the martingale score for + triggering alert. + + :param change_history_length: The length of the sliding window on p-value + for computing the martingale score. + + :param martingale: The type of martingale betting function used for + computing the martingale score. Available options are {``Power``, + ``Mixture``}. + + :param power_martingale_epsilon: The epsilon parameter for the Power + martingale if martingale is set to ``Power``. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: + /../nimbusml/examples/IidSpikeChangePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + confidence=95.0, + change_history_length=20, + martingale='Power', + power_martingale_epsilon=0.1, + **params): + BasePipelineItem.__init__( + self, type='transform', **params) + + self.confidence = confidence + self.change_history_length = change_history_length + self.martingale = martingale + self.power_martingale_epsilon = power_martingale_epsilon + + @property + def _entrypoint(self): + return timeseriesprocessingentrypoints_iidchangepointdetector + + @trace + def _get_node(self, **all_args): + algo_args = dict( + source=self.source, + name=self._name_or_source, + confidence=self.confidence, + change_history_length=self.change_history_length, + martingale=self.martingale, + power_martingale_epsilon=self.power_martingale_epsilon) + + all_args.update(algo_args) + return self._entrypoint(**all_args) diff --git a/src/python/nimbusml/tests/time_series/test_iidchangepointdetector.py b/src/python/nimbusml/tests/time_series/test_iidchangepointdetector.py new file mode 100644 index 00000000..cdaa5691 --- /dev/null +++ b/src/python/nimbusml/tests/time_series/test_iidchangepointdetector.py @@ -0,0 +1,48 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- + +import unittest + +import numpy as np +import pandas as pd +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import IidChangePointDetector + + +class TestIidChangePointDetector(unittest.TestCase): + + def test_correct_data_is_marked_as_change_point(self): + input_data = [5, 5, 5, 5, 5, 5, 5, 5] + input_data.extend([7, 7, 7, 7, 7, 7, 7, 7]) + X_train = pd.Series(input_data, name="ts") + + cpd = IidChangePointDetector(confidence=95, change_history_length=4) << {'result': 'ts'} + data = cpd.fit_transform(X_train) + + self.assertEqual(data.loc[8, 'result.Alert'], 1.0) + + data = data.loc[data['result.Alert'] == 1.0] + self.assertEqual(len(data), 1) + + def test_multiple_user_specified_columns_is_not_allowed(self): + path = get_dataset('timeseries').as_filepath() + data = FileDataStream.read_csv(path) + + try: + pipeline = Pipeline([ + IidChangePointDetector(columns=['t2', 't3'], change_history_length=5) + ]) + pipeline.fit_transform(data) + + except RuntimeError as e: + self.assertTrue('Only one column is allowed' in str(e)) + return + + self.fail() + + +if __name__ == '__main__': + unittest.main() diff --git a/src/python/nimbusml/time_series/__init__.py b/src/python/nimbusml/time_series/__init__.py index 60f1cd18..568e10c3 100644 --- a/src/python/nimbusml/time_series/__init__.py +++ b/src/python/nimbusml/time_series/__init__.py @@ -1,5 +1,7 @@ from .iidspikedetector import IidSpikeDetector +from .iidchangepointdetector import IidChangePointDetector __all__ = [ - 'IidSpikeDetector' + 'IidSpikeDetector', + 'IidChangePointDetector' ] diff --git a/src/python/nimbusml/time_series/iidchangepointdetector.py b/src/python/nimbusml/time_series/iidchangepointdetector.py new file mode 100644 index 00000000..24d6c101 --- /dev/null +++ b/src/python/nimbusml/time_series/iidchangepointdetector.py @@ -0,0 +1,119 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +IidChangePointDetector +""" + +__all__ = ["IidChangePointDetector"] + + +from sklearn.base import TransformerMixin + +from ..base_transform import BaseTransform +from ..internal.core.time_series.iidchangepointdetector import \ + IidChangePointDetector as core +from ..internal.utils.utils import trace + + +class IidChangePointDetector( + core, + BaseTransform, + TransformerMixin): + """ + + This transform detects the change-points in an i.i.d. sequence using + adaptive kernel density estimation and martingales. + + .. remarks:: + ``IIDChangePointDetector`` assumes a sequence of data points that are + independently sampled from one + stationary distribution. `Adaptive kernel density estimation + `_ + is used to model the distribution. + + This transform detects + change points by calculating the martingale score for the sliding + window based on the estimated distribution. + The idea is based on the `Exchangeability + Martingales `_ that + detects a change of distribution over a stream of i.i.d. values. In + short, the value of the + martingale score starts increasing significantly when a sequence of + small p-values are detected in a row; this + indicates the change of the distribution of the underlying data + generation process. + + :param columns: see `Columns `_. + + :param confidence: The confidence for change point detection in the range + [0, 100]. Used to set the threshold of the martingale score for + triggering alert. + + :param change_history_length: The length of the sliding window on p-value + for computing the martingale score. + + :param martingale: The type of martingale betting function used for + computing the martingale score. Available options are {``Power``, + ``Mixture``}. + + :param power_martingale_epsilon: The epsilon parameter for the Power + martingale if martingale is set to ``Power``. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: + /../nimbusml/examples/IidSpikeChangePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + confidence=95.0, + change_history_length=20, + martingale='Power', + power_martingale_epsilon=0.1, + columns=None, + **params): + + if columns: + params['columns'] = columns + BaseTransform.__init__(self, **params) + core.__init__( + self, + confidence=confidence, + change_history_length=change_history_length, + martingale=martingale, + power_martingale_epsilon=power_martingale_epsilon, + **params) + self._columns = columns + + def get_params(self, deep=False): + """ + Get the parameters for this operator. + """ + return core.get_params(self) + + def _nodes_with_presteps(self): + """ + Inserts preprocessing before this one. + """ + from ..preprocessing.schema import TypeConverter + return [ + TypeConverter( + result_type='R4')._steal_io(self), + self] diff --git a/src/python/tests/test_estimator_checks.py b/src/python/tests/test_estimator_checks.py index 1d36dd4d..20006cff 100644 --- a/src/python/tests/test_estimator_checks.py +++ b/src/python/tests/test_estimator_checks.py @@ -16,7 +16,7 @@ from nimbusml.internal.entrypoints._ngramextractor_ngram import n_gram from nimbusml.preprocessing import TensorFlowScorer from nimbusml.preprocessing.filter import SkipFilter, TakeFilter -from nimbusml.time_series import IidSpikeDetector +from nimbusml.time_series import IidSpikeDetector, IidChangePointDetector from sklearn.utils.estimator_checks import _yield_all_checks, MULTI_OUTPUT this = os.path.abspath(os.path.dirname(__file__)) @@ -54,8 +54,9 @@ # fix pending in PR, bug cant handle csr matrix 'RangeFilter': 'check_estimators_dtypes, ' 'check_estimator_sparse_data', - # time series does not currently support sparse matrices + # time series do not currently support sparse matrices 'IidSpikeDetector': 'check_estimator_sparse_data', + 'IidChangePointDetector': 'check_estimator_sparse_data', # bug, low tolerance 'FastLinearRegressor': 'check_supervised_y_2d, ' 'check_regressor_data_not_an_array, ' @@ -184,6 +185,7 @@ 'SkipFilter': SkipFilter(count=5), 'TakeFilter': TakeFilter(count=100000), 'IidSpikeDetector': IidSpikeDetector(columns=['F0']), + 'IidChangePointDetector': IidChangePointDetector(columns=['F0']), 'TensorFlowScorer': TensorFlowScorer( model_location=os.path.join( this, diff --git a/src/python/tools/compiler_utils.py b/src/python/tools/compiler_utils.py index 25faa82b..4039905a 100644 --- a/src/python/tools/compiler_utils.py +++ b/src/python/tools/compiler_utils.py @@ -132,6 +132,7 @@ def _nodes_with_presteps(self): # 'SupervisedBinner': int_to_r4_converter, # not exist in nimbusml 'IidSpikeDetector': timeseries_to_r4_converter, + 'IidChangePointDetector': timeseries_to_r4_converter, 'PcaTransformer': '''from ..preprocessing.schema import TypeConverter diff --git a/src/python/tools/manifest_diff.json b/src/python/tools/manifest_diff.json index 330a49f2..4b3ae241 100644 --- a/src/python/tools/manifest_diff.json +++ b/src/python/tools/manifest_diff.json @@ -545,6 +545,12 @@ "Module": "time_series", "Type": "Transform" }, + { + "Name": "TimeSeriesProcessingEntryPoints.IidChangePointDetector", + "NewName": "IidChangePointDetector", + "Module": "time_series", + "Type": "Transform" + }, { "Name": "Trainers.PoissonRegressor", "NewName": "PoissonRegressionRegressor", From d610dd6a2fb2fccad442522153f41fb44901bfbd Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Fri, 14 Jun 2019 13:42:15 -0700 Subject: [PATCH 10/12] Initial implementation of SsaSpikeDetector. --- src/python/nimbusml.pyproj | 5 + .../nimbusml/examples/SsaSpikeDetector.py | 40 ++++++ .../SsaSpikeDetector_df.py | 80 +++++++++++ .../core/time_series/ssaspikedetector.py | 129 +++++++++++++++++ .../time_series/test_ssaspikedetector.py | 60 ++++++++ src/python/nimbusml/time_series/__init__.py | 4 +- .../nimbusml/time_series/ssaspikedetector.py | 136 ++++++++++++++++++ src/python/tests/test_estimator_checks.py | 4 +- src/python/tools/compiler_utils.py | 1 + src/python/tools/manifest_diff.json | 6 + 10 files changed, 463 insertions(+), 2 deletions(-) create mode 100644 src/python/nimbusml/examples/SsaSpikeDetector.py create mode 100644 src/python/nimbusml/examples/examples_from_dataframe/SsaSpikeDetector_df.py create mode 100644 src/python/nimbusml/internal/core/time_series/ssaspikedetector.py create mode 100644 src/python/nimbusml/tests/time_series/test_ssaspikedetector.py create mode 100644 src/python/nimbusml/time_series/ssaspikedetector.py diff --git a/src/python/nimbusml.pyproj b/src/python/nimbusml.pyproj index e1924f93..e081f6bd 100644 --- a/src/python/nimbusml.pyproj +++ b/src/python/nimbusml.pyproj @@ -89,6 +89,7 @@ + @@ -137,6 +138,7 @@ + @@ -230,6 +232,7 @@ + @@ -579,10 +582,12 @@ + + diff --git a/src/python/nimbusml/examples/SsaSpikeDetector.py b/src/python/nimbusml/examples/SsaSpikeDetector.py new file mode 100644 index 00000000..819f8bc2 --- /dev/null +++ b/src/python/nimbusml/examples/SsaSpikeDetector.py @@ -0,0 +1,40 @@ +############################################################################### +# SsaSpikeDetector +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import SsaSpikeDetector + +# data input (as a FileDataStream) +path = get_dataset('timeseries').as_filepath() + +data = FileDataStream.read_csv(path) +print(data.head()) +# t1 t2 t3 +# 0 0.01 0.01 0.0100 +# 1 0.02 0.02 0.0200 +# 2 0.03 0.03 0.0200 +# 3 0.03 0.03 0.0250 +# 4 0.03 0.03 0.0005 + +# define the training pipeline +pipeline = Pipeline([ + SsaSpikeDetector(columns={'t2_spikes': 't2'}, + pvalue_history_length=4, + training_window_size=8, + seasonal_window_size=3) +]) + +result = pipeline.fit_transform(data) +print(result) + +# t1 t2 t3 t2_spikes.Alert t2_spikes.Raw Score t2_spikes.P-Value Score +# 0 0.01 0.01 0.0100 0.0 -0.111334 5.000000e-01 +# 1 0.02 0.02 0.0200 0.0 -0.076755 4.862075e-01 +# 2 0.03 0.03 0.0200 0.0 -0.034871 3.856320e-03 +# 3 0.03 0.03 0.0250 0.0 -0.012559 8.617091e-02 +# 4 0.03 0.03 0.0005 0.0 -0.015723 2.252377e-01 +# 5 0.03 0.05 0.0100 0.0 -0.001133 1.767711e-01 +# 6 0.05 0.07 0.0500 0.0 0.006265 9.170460e-02 +# 7 0.07 0.09 0.0900 0.0 0.002383 2.701134e-01 +# 8 0.09 99.00 99.0000 1.0 98.879520 1.000000e-08 +# 9 1.10 0.10 0.1000 0.0 -57.817568 6.635692e-02 diff --git a/src/python/nimbusml/examples/examples_from_dataframe/SsaSpikeDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/SsaSpikeDetector_df.py new file mode 100644 index 00000000..0e0196a0 --- /dev/null +++ b/src/python/nimbusml/examples/examples_from_dataframe/SsaSpikeDetector_df.py @@ -0,0 +1,80 @@ +############################################################################### +# SsaSpikeDetector +import numpy as np +import pandas as pd +from nimbusml.time_series import SsaSpikeDetector + +# This example creates a time series (list of data with the +# i-th element corresponding to the i-th time slot). +# The estimator is applied to identify spiking points in the series. +# This estimator can account for temporal seasonality in the data. + +# Generate sample series data with a recurring +# pattern and a spike within the pattern +seasonality_size = 5 +seasonal_data = np.arange(seasonality_size) + +data = np.tile(seasonal_data, 3) +data = np.append(data, [100]) # add a spike +data = np.append(data, seasonal_data) + +X_train = pd.Series(data, name="ts") + +# X_train looks like this +# 0 0 +# 1 1 +# 2 2 +# 3 3 +# 4 4 +# 5 0 +# 6 1 +# 7 2 +# 8 3 +# 9 4 +# 10 0 +# 11 1 +# 12 2 +# 13 3 +# 14 4 +# 15 100 +# 16 0 +# 17 1 +# 18 2 +# 19 3 +# 20 4 + +training_seasons = 3 +training_size = seasonality_size * training_seasons + +ssd = SsaSpikeDetector(confidence=95, + pvalue_history_length=8, + training_window_size=training_size, + seasonal_window_size=seasonality_size + 1) << {'result': 'ts'} + +ssd.fit(X_train, verbose=1) +data = ssd.transform(X_train) + +print(data) + +# ts result.Alert result.Raw Score result.P-Value Score +# 0 0 0.0 -2.531824 5.000000e-01 +# 1 1 0.0 -0.008832 5.818072e-03 +# 2 2 0.0 0.763040 1.374071e-01 +# 3 3 0.0 0.693811 2.797713e-01 +# 4 4 0.0 1.442079 1.838294e-01 +# 5 0 0.0 -1.844414 1.707238e-01 +# 6 1 0.0 0.219578 4.364025e-01 +# 7 2 0.0 0.201708 4.505472e-01 +# 8 3 0.0 0.157089 4.684456e-01 +# 9 4 0.0 1.329494 1.773046e-01 +# 10 0 0.0 -1.792391 7.353794e-02 +# 11 1 0.0 0.161634 4.999295e-01 +# 12 2 0.0 0.092626 4.953789e-01 +# 13 3 0.0 0.084648 4.514174e-01 +# 14 4 0.0 1.305554 1.202619e-01 +# 15 100 1.0 98.207609 1.000000e-08 <-- alert is on, predicted spike +# 16 0 0.0 -13.831450 2.912225e-01 +# 17 1 0.0 -1.741884 4.379857e-01 +# 18 2 0.0 -0.465426 4.557261e-01 +# 19 3 0.0 -16.497133 2.926521e-01 +# 20 4 0.0 -29.817375 2.060473e-01 diff --git a/src/python/nimbusml/internal/core/time_series/ssaspikedetector.py b/src/python/nimbusml/internal/core/time_series/ssaspikedetector.py new file mode 100644 index 00000000..6a1097f8 --- /dev/null +++ b/src/python/nimbusml/internal/core/time_series/ssaspikedetector.py @@ -0,0 +1,129 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +SsaSpikeDetector +""" + +__all__ = ["SsaSpikeDetector"] + + +from ...entrypoints.timeseriesprocessingentrypoints_ssaspikedetector import \ + timeseriesprocessingentrypoints_ssaspikedetector +from ...utils.utils import trace +from ..base_pipeline_item import BasePipelineItem, DefaultSignature + + +class SsaSpikeDetector(BasePipelineItem, DefaultSignature): + """ + + This transform detects the spikes in a seasonal time-series using + Singular Spectrum Analysis (SSA). + + .. remarks:: + `Singular Spectrum Analysis (SSA) + `_ is a + powerful + framework for decomposing the time-series into trend, seasonality and + noise components as well as forecasting + the future values of the time-series. In order to remove the effect + of such components on anomaly detection, + this transform adds SSA as a time-series modeler component in the + detection pipeline. + + The SSA component will be trained and it predicts the next expected + value on the time-series under normal condition; this expected value + is + further used to calculate the amount of deviation from the normal + (predicted) behavior at that timestamp. + The distribution of this deviation is then modeled using `Adaptive + kernel density estimation + `_. + + The `p-value score for the + current deviation is calculated based on the + estimated distribution. The lower its value, the more likely the + current point is an outlier. + + :param training_window_size: The number of points, N, from the beginning + of the sequence used to train the SSA + model. + + :param confidence: The confidence for spike detection in the range [0, + 100]. + + :param seasonal_window_size: An upper bound, L, on the largest relevant + seasonality in the input time-series, which + also determines the order of the autoregression of SSA. It must + satisfy 2 < L < N/2. + + :param side: The argument that determines whether to detect positive or + negative anomalies, or both. Available + options are {``Positive``, ``Negative``, ``TwoSided``}. + + :param pvalue_history_length: The size of the sliding window for computing + the p-value. + + :param error_function: The function used to compute the error between the + expected and the observed value. Possible + values are {``SignedDifference``, ``AbsoluteDifference``, + ``SignedProportion``, ``AbsoluteProportion``, + ``SquaredDifference``}. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/SsaSpikeDetector.py + :language: python + """ + + @trace + def __init__( + self, + training_window_size=100, + confidence=99.0, + seasonal_window_size=10, + side='TwoSided', + pvalue_history_length=100, + error_function='SignedDifference', + **params): + BasePipelineItem.__init__( + self, type='transform', **params) + + self.training_window_size = training_window_size + self.confidence = confidence + self.seasonal_window_size = seasonal_window_size + self.side = side + self.pvalue_history_length = pvalue_history_length + self.error_function = error_function + + @property + def _entrypoint(self): + return timeseriesprocessingentrypoints_ssaspikedetector + + @trace + def _get_node(self, **all_args): + algo_args = dict( + source=self.source, + name=self._name_or_source, + training_window_size=self.training_window_size, + confidence=self.confidence, + seasonal_window_size=self.seasonal_window_size, + side=self.side, + pvalue_history_length=self.pvalue_history_length, + error_function=self.error_function) + + all_args.update(algo_args) + return self._entrypoint(**all_args) diff --git a/src/python/nimbusml/tests/time_series/test_ssaspikedetector.py b/src/python/nimbusml/tests/time_series/test_ssaspikedetector.py new file mode 100644 index 00000000..3645860b --- /dev/null +++ b/src/python/nimbusml/tests/time_series/test_ssaspikedetector.py @@ -0,0 +1,60 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- + +import unittest + +import numpy as np +import pandas as pd +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import SsaSpikeDetector + + +class TestSsaSpikeDetector(unittest.TestCase): + + def test_correct_data_is_marked_as_spike(self): + seasonality_size = 5 + seasonal_data = np.arange(seasonality_size) + + data = np.tile(seasonal_data, 3) + data = np.append(data, [100]) # add a spike + data = np.append(data, seasonal_data) + + X_train = pd.Series(data, name="ts") + training_seasons = 3 + training_size = seasonality_size * training_seasons + + ssd = SsaSpikeDetector(confidence=95, + pvalue_history_length=8, + training_window_size=training_size, + seasonal_window_size=seasonality_size + 1) << {'result': 'ts'} + + ssd.fit(X_train) + data = ssd.transform(X_train) + + self.assertEqual(data.loc[15, 'result.Alert'], 1.0) + + data = data.loc[data['result.Alert'] == 1.0] + self.assertEqual(len(data), 1) + + def test_multiple_user_specified_columns_is_not_allowed(self): + path = get_dataset('timeseries').as_filepath() + data = FileDataStream.read_csv(path) + + try: + pipeline = Pipeline([ + SsaSpikeDetector(columns=['t2', 't3'], pvalue_history_length=5) + ]) + pipeline.fit_transform(data) + + except RuntimeError as e: + self.assertTrue('Only one column is allowed' in str(e)) + return + + self.fail() + + +if __name__ == '__main__': + unittest.main() diff --git a/src/python/nimbusml/time_series/__init__.py b/src/python/nimbusml/time_series/__init__.py index 568e10c3..0f5f72e1 100644 --- a/src/python/nimbusml/time_series/__init__.py +++ b/src/python/nimbusml/time_series/__init__.py @@ -1,7 +1,9 @@ from .iidspikedetector import IidSpikeDetector from .iidchangepointdetector import IidChangePointDetector +from .ssaspikedetector import SsaSpikeDetector __all__ = [ 'IidSpikeDetector', - 'IidChangePointDetector' + 'IidChangePointDetector', + 'SsaSpikeDetector' ] diff --git a/src/python/nimbusml/time_series/ssaspikedetector.py b/src/python/nimbusml/time_series/ssaspikedetector.py new file mode 100644 index 00000000..d57cc4ad --- /dev/null +++ b/src/python/nimbusml/time_series/ssaspikedetector.py @@ -0,0 +1,136 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +SsaSpikeDetector +""" + +__all__ = ["SsaSpikeDetector"] + + +from sklearn.base import TransformerMixin + +from ..base_transform import BaseTransform +from ..internal.core.time_series.ssaspikedetector import \ + SsaSpikeDetector as core +from ..internal.utils.utils import trace + + +class SsaSpikeDetector(core, BaseTransform, TransformerMixin): + """ + + This transform detects the spikes in a seasonal time-series using + Singular Spectrum Analysis (SSA). + + .. remarks:: + `Singular Spectrum Analysis (SSA) + `_ is a + powerful + framework for decomposing the time-series into trend, seasonality and + noise components as well as forecasting + the future values of the time-series. In order to remove the effect + of such components on anomaly detection, + this transform adds SSA as a time-series modeler component in the + detection pipeline. + + The SSA component will be trained and it predicts the next expected + value on the time-series under normal condition; this expected value + is + further used to calculate the amount of deviation from the normal + (predicted) behavior at that timestamp. + The distribution of this deviation is then modeled using `Adaptive + kernel density estimation + `_. + + The `p-value score for the + current deviation is calculated based on the + estimated distribution. The lower its value, the more likely the + current point is an outlier. + + :param columns: see `Columns `_. + + :param training_window_size: The number of points, N, from the beginning + of the sequence used to train the SSA + model. + + :param confidence: The confidence for spike detection in the range [0, + 100]. + + :param seasonal_window_size: An upper bound, L, on the largest relevant + seasonality in the input time-series, which + also determines the order of the autoregression of SSA. It must + satisfy 2 < L < N/2. + + :param side: The argument that determines whether to detect positive or + negative anomalies, or both. Available + options are {``Positive``, ``Negative``, ``TwoSided``}. + + :param pvalue_history_length: The size of the sliding window for computing + the p-value. + + :param error_function: The function used to compute the error between the + expected and the observed value. Possible + values are {``SignedDifference``, ``AbsoluteDifference``, + ``SignedProportion``, ``AbsoluteProportion``, + ``SquaredDifference``}. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaChangePointDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/SsaSpikeDetector.py + :language: python + """ + + @trace + def __init__( + self, + training_window_size=100, + confidence=99.0, + seasonal_window_size=10, + side='TwoSided', + pvalue_history_length=100, + error_function='SignedDifference', + columns=None, + **params): + + if columns: + params['columns'] = columns + BaseTransform.__init__(self, **params) + core.__init__( + self, + training_window_size=training_window_size, + confidence=confidence, + seasonal_window_size=seasonal_window_size, + side=side, + pvalue_history_length=pvalue_history_length, + error_function=error_function, + **params) + self._columns = columns + + def get_params(self, deep=False): + """ + Get the parameters for this operator. + """ + return core.get_params(self) + + def _nodes_with_presteps(self): + """ + Inserts preprocessing before this one. + """ + from ..preprocessing.schema import TypeConverter + return [ + TypeConverter( + result_type='R4')._steal_io(self), + self] diff --git a/src/python/tests/test_estimator_checks.py b/src/python/tests/test_estimator_checks.py index 20006cff..c85e73b6 100644 --- a/src/python/tests/test_estimator_checks.py +++ b/src/python/tests/test_estimator_checks.py @@ -16,7 +16,7 @@ from nimbusml.internal.entrypoints._ngramextractor_ngram import n_gram from nimbusml.preprocessing import TensorFlowScorer from nimbusml.preprocessing.filter import SkipFilter, TakeFilter -from nimbusml.time_series import IidSpikeDetector, IidChangePointDetector +from nimbusml.time_series import IidSpikeDetector, IidChangePointDetector, SsaSpikeDetector from sklearn.utils.estimator_checks import _yield_all_checks, MULTI_OUTPUT this = os.path.abspath(os.path.dirname(__file__)) @@ -57,6 +57,7 @@ # time series do not currently support sparse matrices 'IidSpikeDetector': 'check_estimator_sparse_data', 'IidChangePointDetector': 'check_estimator_sparse_data', + 'SsaSpikeDetector': 'check_estimator_sparse_data', # bug, low tolerance 'FastLinearRegressor': 'check_supervised_y_2d, ' 'check_regressor_data_not_an_array, ' @@ -186,6 +187,7 @@ 'TakeFilter': TakeFilter(count=100000), 'IidSpikeDetector': IidSpikeDetector(columns=['F0']), 'IidChangePointDetector': IidChangePointDetector(columns=['F0']), + 'SsaSpikeDetector': IidChangePointDetector(columns=['F0']), 'TensorFlowScorer': TensorFlowScorer( model_location=os.path.join( this, diff --git a/src/python/tools/compiler_utils.py b/src/python/tools/compiler_utils.py index 4039905a..7e8ef726 100644 --- a/src/python/tools/compiler_utils.py +++ b/src/python/tools/compiler_utils.py @@ -133,6 +133,7 @@ def _nodes_with_presteps(self): 'IidSpikeDetector': timeseries_to_r4_converter, 'IidChangePointDetector': timeseries_to_r4_converter, + 'SsaSpikeDetector': timeseries_to_r4_converter, 'PcaTransformer': '''from ..preprocessing.schema import TypeConverter diff --git a/src/python/tools/manifest_diff.json b/src/python/tools/manifest_diff.json index 4b3ae241..696949a9 100644 --- a/src/python/tools/manifest_diff.json +++ b/src/python/tools/manifest_diff.json @@ -551,6 +551,12 @@ "Module": "time_series", "Type": "Transform" }, + { + "Name": "TimeSeriesProcessingEntryPoints.SsaSpikeDetector", + "NewName": "SsaSpikeDetector", + "Module": "time_series", + "Type": "Transform" + }, { "Name": "Trainers.PoissonRegressor", "NewName": "PoissonRegressionRegressor", From 97dad7a80805aa32ec215fa65da6d306a3a3d0c9 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Fri, 14 Jun 2019 15:27:15 -0700 Subject: [PATCH 11/12] Initial implementation of SsaChangePointDetector. --- src/python/nimbusml.pyproj | 5 + .../examples/SsaChangePointDetector.py | 40 +++++ .../SsaChangePointDetector_df.py | 77 +++++++++ .../time_series/ssachangepointdetector.py | 138 ++++++++++++++++ .../test_ssachangepointdetector.py | 61 ++++++++ src/python/nimbusml/time_series/__init__.py | 4 +- .../time_series/ssachangepointdetector.py | 147 ++++++++++++++++++ src/python/tests/test_estimator_checks.py | 6 +- src/python/tools/compiler_utils.py | 1 + src/python/tools/manifest_diff.json | 6 + 10 files changed, 483 insertions(+), 2 deletions(-) create mode 100644 src/python/nimbusml/examples/SsaChangePointDetector.py create mode 100644 src/python/nimbusml/examples/examples_from_dataframe/SsaChangePointDetector_df.py create mode 100644 src/python/nimbusml/internal/core/time_series/ssachangepointdetector.py create mode 100644 src/python/nimbusml/tests/time_series/test_ssachangepointdetector.py create mode 100644 src/python/nimbusml/time_series/ssachangepointdetector.py diff --git a/src/python/nimbusml.pyproj b/src/python/nimbusml.pyproj index e081f6bd..9c09758d 100644 --- a/src/python/nimbusml.pyproj +++ b/src/python/nimbusml.pyproj @@ -89,6 +89,7 @@ + @@ -138,6 +139,7 @@ + @@ -232,6 +234,7 @@ + @@ -582,11 +585,13 @@ + + diff --git a/src/python/nimbusml/examples/SsaChangePointDetector.py b/src/python/nimbusml/examples/SsaChangePointDetector.py new file mode 100644 index 00000000..e797bc30 --- /dev/null +++ b/src/python/nimbusml/examples/SsaChangePointDetector.py @@ -0,0 +1,40 @@ +############################################################################### +# SsaChangePointDetector +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import SsaChangePointDetector + +# data input (as a FileDataStream) +path = get_dataset('timeseries').as_filepath() + +data = FileDataStream.read_csv(path) +print(data.head()) +# t1 t2 t3 +# 0 0.01 0.01 0.0100 +# 1 0.02 0.02 0.0200 +# 2 0.03 0.03 0.0200 +# 3 0.03 0.03 0.0250 +# 4 0.03 0.03 0.0005 + +# define the training pipeline +pipeline = Pipeline([ + SsaChangePointDetector(columns={'t2_cp': 't2'}, + change_history_length=4, + training_window_size=8, + seasonal_window_size=3) +]) + +result = pipeline.fit_transform(data) +print(result) + +# t1 t2 t3 t2_cp.Alert t2_cp.Raw Score t2_cp.P-Value Score t2_cp.Martingale Score +# 0 0.01 0.01 0.0100 0.0 -0.111334 5.000000e-01 0.001213 +# 1 0.02 0.02 0.0200 0.0 -0.076755 4.862075e-01 0.001243 +# 2 0.03 0.03 0.0200 0.0 -0.034871 3.856320e-03 0.099119 +# 3 0.03 0.03 0.0250 0.0 -0.012559 8.617091e-02 0.482400 +# 4 0.03 0.03 0.0005 0.0 -0.015723 2.252377e-01 0.988788 +# 5 0.03 0.05 0.0100 0.0 -0.001133 1.767711e-01 2.457946 +# 6 0.05 0.07 0.0500 0.0 0.006265 9.170460e-02 0.141898 +# 7 0.07 0.09 0.0900 0.0 0.002383 2.701134e-01 0.050747 +# 8 0.09 99.00 99.0000 1.0 98.879520 1.000000e-08 210274.372059 +# 9 1.10 0.10 0.1000 0.0 -57.817568 6.635692e-02 507877.454862 diff --git a/src/python/nimbusml/examples/examples_from_dataframe/SsaChangePointDetector_df.py b/src/python/nimbusml/examples/examples_from_dataframe/SsaChangePointDetector_df.py new file mode 100644 index 00000000..8f1a027d --- /dev/null +++ b/src/python/nimbusml/examples/examples_from_dataframe/SsaChangePointDetector_df.py @@ -0,0 +1,77 @@ +############################################################################### +# SsaChangePointDetector +import numpy as np +import pandas as pd +from nimbusml.time_series import SsaChangePointDetector + +# This example creates a time series (list of data with the +# i-th element corresponding to the i-th time slot). +# The estimator is applied to identify points where data distribution changed. +# This estimator can account for temporal seasonality in the data. + +# Generate sample series data with a recurring +# pattern and a spike within the pattern +seasonality_size = 5 +seasonal_data = np.arange(seasonality_size) + +data = np.tile(seasonal_data, 3) +data = np.append(data, [0, 100, 200, 300, 400]) # change distribution + +X_train = pd.Series(data, name="ts") + +# X_train looks like this +# 0 0 +# 1 1 +# 2 2 +# 3 3 +# 4 4 +# 5 0 +# 6 1 +# 7 2 +# 8 3 +# 9 4 +# 10 0 +# 11 1 +# 12 2 +# 13 3 +# 14 4 +# 15 0 +# 16 100 +# 17 200 +# 18 300 +# 19 400 + +training_seasons = 3 +training_size = seasonality_size * training_seasons + +cpd = SsaChangePointDetector(confidence=95, + change_history_length=8, + training_window_size=training_size, + seasonal_window_size=seasonality_size + 1) << {'result': 'ts'} + +cpd.fit(X_train, verbose=1) +data = cpd.transform(X_train) + +print(data) + +# ts result.Alert result.Raw Score result.P-Value Score result.Martingale Score +# 0 0 0.0 -2.531824 5.000000e-01 1.470334e-06 +# 1 1 0.0 -0.008832 5.818072e-03 8.094459e-05 +# 2 2 0.0 0.763040 1.374071e-01 2.588526e-04 +# 3 3 0.0 0.693811 2.797713e-01 4.365186e-04 +# 4 4 0.0 1.442079 1.838294e-01 1.074242e-03 +# 5 0 0.0 -1.844414 1.707238e-01 2.825599e-03 +# 6 1 0.0 0.219578 4.364025e-01 3.193633e-03 +# 7 2 0.0 0.201708 4.505472e-01 3.507451e-03 +# 8 3 0.0 0.157089 4.684456e-01 3.719387e-03 +# 9 4 0.0 1.329494 1.773046e-01 1.717610e-04 +# 10 0 0.0 -1.792391 7.353794e-02 3.014897e-04 +# 11 1 0.0 0.161634 4.999295e-01 1.788041e-04 +# 12 2 0.0 0.092626 4.953789e-01 7.326680e-05 +# 13 3 0.0 0.084648 4.514174e-01 3.053876e-05 +# 14 4 0.0 1.305554 1.202619e-01 9.741702e-05 +# 15 0 0.0 -1.792391 7.264402e-02 5.034093e-04 +# 16 100 1.0 99.161634 1.000000e-08 4.031944e+03 <-- alert is on, predicted spike +# 17 200 0.0 185.229474 5.485437e-04 7.312609e+05 +# 18 300 0.0 270.403543 1.259683e-02 3.578470e+06 +# 19 400 0.0 357.113747 2.978766e-02 4.529837e+07 diff --git a/src/python/nimbusml/internal/core/time_series/ssachangepointdetector.py b/src/python/nimbusml/internal/core/time_series/ssachangepointdetector.py new file mode 100644 index 00000000..297fae42 --- /dev/null +++ b/src/python/nimbusml/internal/core/time_series/ssachangepointdetector.py @@ -0,0 +1,138 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +SsaChangePointDetector +""" + +__all__ = ["SsaChangePointDetector"] + + +from ...entrypoints.timeseriesprocessingentrypoints_ssachangepointdetector import \ + timeseriesprocessingentrypoints_ssachangepointdetector +from ...utils.utils import trace +from ..base_pipeline_item import BasePipelineItem, DefaultSignature + + +class SsaChangePointDetector(BasePipelineItem, DefaultSignature): + """ + + This transform detects the change-points in a seasonal time-series + using Singular Spectrum Analysis (SSA). + + .. remarks:: + `Singular Spectrum Analysis (SSA) + `_ is a + powerful framework for decomposing the time-series into trend, + seasonality and noise components as well as forecasting the future + values of the time-series. In order to remove the + effect of such components on anomaly detection, this transform add + SSA as a time-series modeler component in the detection pipeline. + + The SSA component will be trained and it predicts the next expected + value on the time-series under normal condition; this expected value + is + further used to calculate the amount of deviation from the normal + behavior at that timestamp. + The distribution of this deviation is then modeled using `Adaptive + kernel density estimation + `_. + + This transform detects + change points by calculating the martingale score for the sliding + window based on the estimated distribution of deviations. + The idea is based on the `Exchangeability + Martingales `_ that + detects a change of distribution over a stream of i.i.d. values. In + short, the value of the + martingale score starts increasing significantly when a sequence of + small p-values detected in a row; this + indicates the change of the distribution of the underlying data + generation process. + + :param training_window_size: The number of points, N, from the beginning + of the sequence used to train the SSA model. + + :param confidence: The confidence for change point detection in the range + [0, 100]. + + :param seasonal_window_size: An upper bound, L, on the largest relevant + seasonality in the input time-series, which also + determines the order of the autoregression of SSA. It must satisfy 2 + < L < N/2. + + :param change_history_length: The length of the sliding window on p-value + for computing the martingale score. + + :param error_function: The function used to compute the error between the + expected and the observed value. Possible values are: + {``SignedDifference``, ``AbsoluteDifference``, ``SignedProportion``, + ``AbsoluteProportion``, ``SquaredDifference``}. + + :param martingale: The type of martingale betting function used for + computing the martingale score. Available options are {``Power``, + ``Mixture``}. + + :param power_martingale_epsilon: The epsilon parameter for the Power + martingale if martingale is set to ``Power``. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaSpikeDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/SsaChangePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + training_window_size=100, + confidence=95.0, + seasonal_window_size=10, + change_history_length=20, + error_function='SignedDifference', + martingale='Power', + power_martingale_epsilon=0.1, + **params): + BasePipelineItem.__init__( + self, type='transform', **params) + + self.training_window_size = training_window_size + self.confidence = confidence + self.seasonal_window_size = seasonal_window_size + self.change_history_length = change_history_length + self.error_function = error_function + self.martingale = martingale + self.power_martingale_epsilon = power_martingale_epsilon + + @property + def _entrypoint(self): + return timeseriesprocessingentrypoints_ssachangepointdetector + + @trace + def _get_node(self, **all_args): + algo_args = dict( + source=self.source, + name=self._name_or_source, + training_window_size=self.training_window_size, + confidence=self.confidence, + seasonal_window_size=self.seasonal_window_size, + change_history_length=self.change_history_length, + error_function=self.error_function, + martingale=self.martingale, + power_martingale_epsilon=self.power_martingale_epsilon) + + all_args.update(algo_args) + return self._entrypoint(**all_args) diff --git a/src/python/nimbusml/tests/time_series/test_ssachangepointdetector.py b/src/python/nimbusml/tests/time_series/test_ssachangepointdetector.py new file mode 100644 index 00000000..b115396a --- /dev/null +++ b/src/python/nimbusml/tests/time_series/test_ssachangepointdetector.py @@ -0,0 +1,61 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- + +import unittest + +import numpy as np +import pandas as pd +from nimbusml import Pipeline, FileDataStream +from nimbusml.datasets import get_dataset +from nimbusml.time_series import SsaChangePointDetector + + +class TestSsaChangePointDetector(unittest.TestCase): + + def test_correct_data_is_marked_as_change_point(self): + seasonality_size = 5 + seasonal_data = np.arange(seasonality_size) + + data = np.tile(seasonal_data, 3) + data = np.append(data, [0, 100, 200, 300, 400]) # change distribution + + X_train = pd.Series(data, name="ts") + + training_seasons = 3 + training_size = seasonality_size * training_seasons + + cpd = SsaChangePointDetector(confidence=95, + change_history_length=8, + training_window_size=training_size, + seasonal_window_size=seasonality_size + 1) << {'result': 'ts'} + + cpd.fit(X_train, verbose=1) + data = cpd.transform(X_train) + + + self.assertEqual(data.loc[16, 'result.Alert'], 1.0) + + data = data.loc[data['result.Alert'] == 1.0] + self.assertEqual(len(data), 1) + + def test_multiple_user_specified_columns_is_not_allowed(self): + path = get_dataset('timeseries').as_filepath() + data = FileDataStream.read_csv(path) + + try: + pipeline = Pipeline([ + SsaChangePointDetector(columns=['t2', 't3'], change_history_length=5) + ]) + pipeline.fit_transform(data) + + except RuntimeError as e: + self.assertTrue('Only one column is allowed' in str(e)) + return + + self.fail() + + +if __name__ == '__main__': + unittest.main() diff --git a/src/python/nimbusml/time_series/__init__.py b/src/python/nimbusml/time_series/__init__.py index 0f5f72e1..807e3a7b 100644 --- a/src/python/nimbusml/time_series/__init__.py +++ b/src/python/nimbusml/time_series/__init__.py @@ -1,9 +1,11 @@ from .iidspikedetector import IidSpikeDetector from .iidchangepointdetector import IidChangePointDetector from .ssaspikedetector import SsaSpikeDetector +from .ssachangepointdetector import SsaChangePointDetector __all__ = [ 'IidSpikeDetector', 'IidChangePointDetector', - 'SsaSpikeDetector' + 'SsaSpikeDetector', + 'SsaChangePointDetector' ] diff --git a/src/python/nimbusml/time_series/ssachangepointdetector.py b/src/python/nimbusml/time_series/ssachangepointdetector.py new file mode 100644 index 00000000..adf7f9a8 --- /dev/null +++ b/src/python/nimbusml/time_series/ssachangepointdetector.py @@ -0,0 +1,147 @@ +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------------------------- +# - Generated by tools/entrypoint_compiler.py: do not edit by hand +""" +SsaChangePointDetector +""" + +__all__ = ["SsaChangePointDetector"] + + +from sklearn.base import TransformerMixin + +from ..base_transform import BaseTransform +from ..internal.core.time_series.ssachangepointdetector import \ + SsaChangePointDetector as core +from ..internal.utils.utils import trace + + +class SsaChangePointDetector( + core, + BaseTransform, + TransformerMixin): + """ + + This transform detects the change-points in a seasonal time-series + using Singular Spectrum Analysis (SSA). + + .. remarks:: + `Singular Spectrum Analysis (SSA) + `_ is a + powerful framework for decomposing the time-series into trend, + seasonality and noise components as well as forecasting the future + values of the time-series. In order to remove the + effect of such components on anomaly detection, this transform add + SSA as a time-series modeler component in the detection pipeline. + + The SSA component will be trained and it predicts the next expected + value on the time-series under normal condition; this expected value + is + further used to calculate the amount of deviation from the normal + behavior at that timestamp. + The distribution of this deviation is then modeled using `Adaptive + kernel density estimation + `_. + + This transform detects + change points by calculating the martingale score for the sliding + window based on the estimated distribution of deviations. + The idea is based on the `Exchangeability + Martingales `_ that + detects a change of distribution over a stream of i.i.d. values. In + short, the value of the + martingale score starts increasing significantly when a sequence of + small p-values detected in a row; this + indicates the change of the distribution of the underlying data + generation process. + + :param columns: see `Columns `_. + + :param training_window_size: The number of points, N, from the beginning + of the sequence used to train the SSA model. + + :param confidence: The confidence for change point detection in the range + [0, 100]. + + :param seasonal_window_size: An upper bound, L, on the largest relevant + seasonality in the input time-series, which also + determines the order of the autoregression of SSA. It must satisfy 2 + < L < N/2. + + :param change_history_length: The length of the sliding window on p-value + for computing the martingale score. + + :param error_function: The function used to compute the error between the + expected and the observed value. Possible values are: + {``SignedDifference``, ``AbsoluteDifference``, ``SignedProportion``, + ``AbsoluteProportion``, ``SquaredDifference``}. + + :param martingale: The type of martingale betting function used for + computing the martingale score. Available options are {``Power``, + ``Mixture``}. + + :param power_martingale_epsilon: The epsilon parameter for the Power + martingale if martingale is set to ``Power``. + + :param params: Additional arguments sent to compute engine. + + .. seealso:: + :py:func:`IIDChangePointDetector + `, + :py:func:`IIDSpikeDetector + `, + :py:func:`SsaSpikeDetector + `. + + .. index:: models, timeseries, transform + + Example: + .. literalinclude:: /../nimbusml/examples/SsaChangePointDetector.py + :language: python + """ + + @trace + def __init__( + self, + training_window_size=100, + confidence=95.0, + seasonal_window_size=10, + change_history_length=20, + error_function='SignedDifference', + martingale='Power', + power_martingale_epsilon=0.1, + columns=None, + **params): + + if columns: + params['columns'] = columns + BaseTransform.__init__(self, **params) + core.__init__( + self, + training_window_size=training_window_size, + confidence=confidence, + seasonal_window_size=seasonal_window_size, + change_history_length=change_history_length, + error_function=error_function, + martingale=martingale, + power_martingale_epsilon=power_martingale_epsilon, + **params) + self._columns = columns + + def get_params(self, deep=False): + """ + Get the parameters for this operator. + """ + return core.get_params(self) + + def _nodes_with_presteps(self): + """ + Inserts preprocessing before this one. + """ + from ..preprocessing.schema import TypeConverter + return [ + TypeConverter( + result_type='R4')._steal_io(self), + self] diff --git a/src/python/tests/test_estimator_checks.py b/src/python/tests/test_estimator_checks.py index c85e73b6..d502b3ce 100644 --- a/src/python/tests/test_estimator_checks.py +++ b/src/python/tests/test_estimator_checks.py @@ -16,7 +16,8 @@ from nimbusml.internal.entrypoints._ngramextractor_ngram import n_gram from nimbusml.preprocessing import TensorFlowScorer from nimbusml.preprocessing.filter import SkipFilter, TakeFilter -from nimbusml.time_series import IidSpikeDetector, IidChangePointDetector, SsaSpikeDetector +from nimbusml.time_series import (IidSpikeDetector, IidChangePointDetector, + SsaSpikeDetector, SsaChangePointDetector) from sklearn.utils.estimator_checks import _yield_all_checks, MULTI_OUTPUT this = os.path.abspath(os.path.dirname(__file__)) @@ -58,6 +59,8 @@ 'IidSpikeDetector': 'check_estimator_sparse_data', 'IidChangePointDetector': 'check_estimator_sparse_data', 'SsaSpikeDetector': 'check_estimator_sparse_data', + 'SsaChangePointDetector': 'check_estimator_sparse_data' + 'check_fit2d_1sample', # SSA requires more than one sample # bug, low tolerance 'FastLinearRegressor': 'check_supervised_y_2d, ' 'check_regressor_data_not_an_array, ' @@ -188,6 +191,7 @@ 'IidSpikeDetector': IidSpikeDetector(columns=['F0']), 'IidChangePointDetector': IidChangePointDetector(columns=['F0']), 'SsaSpikeDetector': IidChangePointDetector(columns=['F0']), + 'SsaChangePointDetector': SsaChangePointDetector(columns=['F0'], seasonal_window_size=2), 'TensorFlowScorer': TensorFlowScorer( model_location=os.path.join( this, diff --git a/src/python/tools/compiler_utils.py b/src/python/tools/compiler_utils.py index 7e8ef726..9a5e1e07 100644 --- a/src/python/tools/compiler_utils.py +++ b/src/python/tools/compiler_utils.py @@ -134,6 +134,7 @@ def _nodes_with_presteps(self): 'IidSpikeDetector': timeseries_to_r4_converter, 'IidChangePointDetector': timeseries_to_r4_converter, 'SsaSpikeDetector': timeseries_to_r4_converter, + 'SsaChangePointDetector': timeseries_to_r4_converter, 'PcaTransformer': '''from ..preprocessing.schema import TypeConverter diff --git a/src/python/tools/manifest_diff.json b/src/python/tools/manifest_diff.json index 696949a9..6c96eb5c 100644 --- a/src/python/tools/manifest_diff.json +++ b/src/python/tools/manifest_diff.json @@ -557,6 +557,12 @@ "Module": "time_series", "Type": "Transform" }, + { + "Name": "TimeSeriesProcessingEntryPoints.SsaChangePointDetector", + "NewName": "SsaChangePointDetector", + "Module": "time_series", + "Type": "Transform" + }, { "Name": "Trainers.PoissonRegressor", "NewName": "PoissonRegressionRegressor", From 1ffec25a533ca1d6e8e7badc77d3f7297887e750 Mon Sep 17 00:00:00 2001 From: "pieths.dev@gmail.com" Date: Fri, 14 Jun 2019 15:36:04 -0700 Subject: [PATCH 12/12] Fix incorrect SsaSpikeDetector instance in test_estimator_checks. --- src/python/tests/test_estimator_checks.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/python/tests/test_estimator_checks.py b/src/python/tests/test_estimator_checks.py index d502b3ce..cf1a358b 100644 --- a/src/python/tests/test_estimator_checks.py +++ b/src/python/tests/test_estimator_checks.py @@ -58,7 +58,8 @@ # time series do not currently support sparse matrices 'IidSpikeDetector': 'check_estimator_sparse_data', 'IidChangePointDetector': 'check_estimator_sparse_data', - 'SsaSpikeDetector': 'check_estimator_sparse_data', + 'SsaSpikeDetector': 'check_estimator_sparse_data' + 'check_fit2d_1sample', # SSA requires more than one sample 'SsaChangePointDetector': 'check_estimator_sparse_data' 'check_fit2d_1sample', # SSA requires more than one sample # bug, low tolerance @@ -190,7 +191,7 @@ 'TakeFilter': TakeFilter(count=100000), 'IidSpikeDetector': IidSpikeDetector(columns=['F0']), 'IidChangePointDetector': IidChangePointDetector(columns=['F0']), - 'SsaSpikeDetector': IidChangePointDetector(columns=['F0']), + 'SsaSpikeDetector': SsaSpikeDetector(columns=['F0'], seasonal_window_size=2), 'SsaChangePointDetector': SsaChangePointDetector(columns=['F0'], seasonal_window_size=2), 'TensorFlowScorer': TensorFlowScorer( model_location=os.path.join(