This is an automated email from the ASF dual-hosted git repository.
claudevdm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new a7b702ed5ba [Python] Deflake anomaly transform tests (#39669)
a7b702ed5ba is described below
commit a7b702ed5ba23973fcbcb529ffb4fc9510aadac1
Author: Bruno Volpato <[email protected]>
AuthorDate: Thu Sep 10 09:27:09 2026 -0400
[Python] Deflake anomaly transform tests (#39669)
* [Python] Deflake anomaly transform tests
* Test ensemble model state without input ordering assumptions
---
.../apache_beam/ml/anomaly/transforms_test.py | 81 ++++++++++++++++++----
1 file changed, 67 insertions(+), 14 deletions(-)
diff --git a/sdks/python/apache_beam/ml/anomaly/transforms_test.py
b/sdks/python/apache_beam/ml/anomaly/transforms_test.py
index 423e51abf63..72c16841d5f 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 @@ from parameterized import parameterized
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 @@ class TestAnomalyDetection(unittest.TestCase):
])
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 @@ class TestAnomalyDetection(unittest.TestCase):
])
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"))
@@ -286,6 +310,35 @@ class TestAnomalyDetection(unittest.TestCase):
else:
assert_that(result, equal_to(expected, _unkeyed_result_is_equal_to))
+ @parameterized.expand([(False, ), (True, )])
+ def test_ensemble_preserves_model_state(self, keyed):
+ # Constant values make the scores independent of processing order. Each
+ # model must produce exactly two warmup predictions per key.
+ input = [beam.Row(x1=1, x2=4)] * 6
+ counts_by_key = [(0, 6)]
+ if keyed:
+ input = [(1, row) for row in input] + [(2, beam.Row(x1=100, x2=50))] * 3
+ counts_by_key = [(1, 6), (2, 3)]
+
+ detectors = [
+ ZScore(features=[feature], model_id=feature)
+ for feature in ('x1', 'x2')
+ ]
+ expected = [(key, model_id, label) for key, count in counts_by_key
+ for model_id in ('x1', 'x2')
+ for label in [-2, -2] + [0] * (count - 2)]
+
+ with TestPipeline() as p:
+ result = (
+ p | beam.Create(input)
+ | AnomalyDetection(EnsembleAnomalyDetector(detectors)))
+ if not keyed:
+ result = result | beam.WithKeys(0)
+ labels = result | beam.FlatMap(
+ lambda item: [(item[0], prediction.model_id, prediction.label)
+ for prediction in item[1].predictions])
+ assert_that(labels, equal_to(expected))
+
class FakeNumpyModel():
def __init__(self):