diff --git a/sdks/python/apache_beam/ml/anomaly/transforms_test.py b/sdks/python/apache_beam/ml/anomaly/transforms_test.py index 423e51abf635..70c00635b1cd 100644 --- a/sdks/python/apache_beam/ml/anomaly/transforms_test.py +++ b/sdks/python/apache_beam/ml/anomaly/transforms_test.py @@ -32,6 +32,7 @@ import apache_beam as beam from apache_beam.ml.anomaly.aggregations import AnyVote +from apache_beam.ml.anomaly.base import AnomalyDetector from apache_beam.ml.anomaly.base import AnomalyPrediction from apache_beam.ml.anomaly.base import AnomalyResult from apache_beam.ml.anomaly.base import EnsembleAnomalyDetector @@ -116,18 +117,33 @@ def _keyed_result_is_equal_to( return a[0] == b[0] and _unkeyed_result_is_equal_to(a[1], b[1]) +@specifiable +class _ValueDetector(AnomalyDetector): + def learn_one(self, unused_x: beam.Row) -> None: + pass + + def score_one(self, x: beam.Row) -> float: + value = float(next(iter(x))) + return float('NaN') if value == -1 else value + + class TestAnomalyDetection(unittest.TestCase): class TestData: unkeyed_input = [ - beam.Row(x1=1, x2=4), - beam.Row(x1=2, x2=4), - beam.Row(x1=3, x2=5), - beam.Row(x1=10, x2=4), # outlier in key=1, with respect to x1 - beam.Row(x1=2, x2=10), # outlier in key=1, with respect to x2 - beam.Row(x1=3, x2=4), + beam.Row(x1=1, x2=4, score_x1=-1.0, score_x2=-1.0), + beam.Row(x1=2, x2=4, score_x1=-1.0, score_x2=-1.0), + beam.Row(x1=3, x2=5, score_x1=2.1213203435596424, score_x2=0.0), + beam.Row(x1=10, x2=4, score_x1=8.0, score_x2=0.5773502691896252), + beam.Row(x1=2, x2=10, score_x1=0.4898979485566356, score_x2=11.5), + beam.Row( + x1=3, + x2=4, + score_x1=0.16452254913212455, + score_x2=0.5368754921931594), + ] + keyed_input = list(zip(itertools.repeat(1), unkeyed_input)) + [ + (2, beam.Row(x1=100, x2=5, score_x1=-1.0, score_x2=-1.0)) ] - keyed_input = list(zip(itertools.repeat(1), - unkeyed_input)) + [(2, beam.Row(x1=100, x2=5))] zscore_x1_expected_predictions = [ AnomalyPrediction( @@ -244,10 +260,14 @@ def test_one_detector(self, input, expected): ]) def test_multiple_detectors_without_aggregation(self, input, expected): sub_detectors = [] - sub_detectors.append(ZScore(features=["x1"], model_id="zscore_x1")) sub_detectors.append( - ZScore( - features=["x2"], + _ValueDetector( + features=["score_x1"], + threshold_criterion=FixedThreshold(3), + model_id="zscore_x1")) + sub_detectors.append( + _ValueDetector( + features=["score_x2"], threshold_criterion=FixedThreshold(2), model_id="zscore_x2")) @@ -267,10 +287,14 @@ def test_multiple_detectors_without_aggregation(self, input, expected): ]) def test_multiple_sub_detectors_with_aggregation(self, input, expected): sub_detectors = [] - sub_detectors.append(ZScore(features=["x1"], model_id="zscore_x1")) sub_detectors.append( - ZScore( - features=["x2"], + _ValueDetector( + features=["score_x1"], + threshold_criterion=FixedThreshold(3), + model_id="zscore_x1")) + sub_detectors.append( + _ValueDetector( + features=["score_x2"], threshold_criterion=FixedThreshold(2), model_id="zscore_x2"))