Multivariate Anomalie-Erkennung mit Isolationswald

Dieser Artikel zeigt, wie man SynapseML auf Apache Spark zur Erkennung multivariater Anomalien verwendet. Die Multivariate Anomalie-Erkennung erkennt Anomalien zwischen vielen Variablen oder Zeitreihen und berücksichtigt dabei alle Wechselwirkungen und Abhängigkeiten zwischen den verschiedenen Variablen. In diesem Szenario nutzt du SynapseML, um ein Isolationswaldmodell für die Erkennung multivariater Anomalien zu trainieren, und nutzt dann das trainierte Modell, um multivariate Anomalien in einem Datensatz mit synthetischen Messungen von drei IoT-Sensoren abzuleiten.

Um mehr über das Isolationswaldmodell zu erfahren, siehe die Originalarbeit von Liu et al.

Voraussetzungen

  1. Erwerben Sie ein Microsoft Fabric-Abonnement. Oder registrieren Sie sich für eine kostenlose Microsoft Fabric-Testversion.
  2. Fügen Sie Ihr Notizbuch an ein Seehaus an. Wählen Sie auf der linken Seite Hinzufügen aus, um ein vorhandenes Seehaus hinzuzufügen oder ein Seehaus zu erstellen.
  3. SynapseML ist in Fabric PySpark-Laufzeiten vorinstalliert (Laufzeit 1.3 oder später wird empfohlen). Um eine bestimmte Version zu verwenden, siehe Install a different version of SynapseML on Fabric.

Bibliotheksimporte

from pyspark.sql import functions as F
from pyspark.ml.feature import VectorAssembler
from pyspark.sql.types import DoubleType
from pyspark.ml import Pipeline

from synapse.ml.isolationforest import IsolationForest
from pyspark.sql import SparkSession

# Bootstrap Spark Session
spark = SparkSession.builder.getOrCreate()

Eingangsdaten

# Table inputs
timestampColumn = "timestamp"  # str: the name of the timestamp column in the table
inputCols = [
    "sensor_1",
    "sensor_2",
    "sensor_3",
]  # list(str): the names of the input variables

# Training Start time, and number of days to use for training:
trainingStartTime = (
    "2022-02-24T06:00:00Z"  # datetime: datetime for when to start the training
)
trainingEndTime = (
    "2022-03-08T23:55:00Z"  # datetime: datetime for when to end the training
)
inferenceStartTime = (
    "2022-03-09T09:30:00Z"  # datetime: datetime for when to start the inference
)
inferenceEndTime = (
    "2022-03-20T23:55:00Z"  # datetime: datetime for when to end the inference
)

# Isolation Forest parameters
contamination = 0.021
num_estimators = 100
max_samples = 256
max_features = 1.0

Daten lesen

df = (
    spark.read.format("csv")
    .option("header", "true")
    .load(
        "wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
    )
)

Konvertieren Sie die Spalten in die entsprechenden Datentypen.

df = (
    df.orderBy(timestampColumn)
    .withColumn("timestamp", F.date_format(timestampColumn, "yyyy-MM-dd'T'HH:mm:ss'Z'"))
    .withColumn("sensor_1", F.col("sensor_1").cast(DoubleType()))
    .withColumn("sensor_2", F.col("sensor_2").cast(DoubleType()))
    .withColumn("sensor_3", F.col("sensor_3").cast(DoubleType()))
    .drop("_c5")  # drop the extra unlabeled column present in the source CSV
)

display(df)

Vorbereitung von Schulungsdaten

# filter to data with timestamps within the training window
df_train = df.filter(
    (F.col(timestampColumn) >= trainingStartTime)
    & (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)

Testen der Datenvorbereitung

# filter to data with timestamps within the inference window
df_test = df.filter(
    (F.col(timestampColumn) >= inferenceStartTime)
    & (F.col(timestampColumn) <= inferenceEndTime)
)
display(df_test)

Isolation-Forest-Modell trainieren

isolationForest = (
    IsolationForest()
    .setNumEstimators(num_estimators)
    .setBootstrap(False)
    .setMaxSamples(max_samples)
    .setMaxFeatures(max_features)
    .setFeaturesCol("features")
    .setPredictionCol("predictedLabel")
    .setScoreCol("outlierScore")
    .setContamination(contamination)
    .setContaminationError(0.01 * contamination)
    .setRandomSeed(1)
)

Erstellen Sie anschließend eine ML-Pipeline, um das Isolation Forest-Modell zu trainieren.

Für das Training des Modells und die Durchführung von Inferenzen im selben Notizbuch reicht das Modellobjekt aus. Um das Modell über Sitzungen hinweg zu erhalten und wiederzuverwenden, registrieren Sie es bei MLflow in Microsoft Fabric.

va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)

Durchführen von Inferenzen

Wenden Sie das trainierte Modell auf die Testdaten an:

df_test_pred = model.transform(df_test)
display(df_test_pred)

Vorgefertigter Anomaliedetektor

Von Bedeutung

Microsoft stellt den Azure KI Anomalie Detektor-Dienst am 1. Oktober 2026 zurück. Seit dem 20. September 2023 können Sie keine neuen Ressourcen mehr erstellen. Für eine unterstützte Alternative, siehe Anomalieerkennung in Microsoft Fabric Real-Time Intelligence.

Azure AI-Anomalie-Detektor

  • Anomaliestatus des letzten Punktes: erzeugt ein Modell unter Verwendung der vorhergehenden Punkte und bestimmt, ob der letzte Punkt anomal ist. Siehe das SynapseML GitHub-Repository für aktuelle API-Referenzen.
  • Anomalien finden: Erzeugt ein Modell mit einer ganzen Serie und findet Anomalien in der Serie. Siehe das SynapseML GitHub-Repository für aktuelle API-Referenzen.