Détection d’anomalies multivariées avec forêt d’isolement

Cet article montre comment utiliser SynapseML sur Apache Spark pour la détection d’anomalies multivariées. La détection d’anomalies multivariées détecte des anomalies parmi de nombreuses variables ou séries temporelles, en tenant compte de toutes les corrélations et dépendances entre les différentes variables. Dans ce scénario, vous utilisez SynapseML pour entraîner un modèle de forêt d’isolation pour la détection d’anomalies multivariées, puis vous utilisez le modèle entraîné pour inférer des anomalies multivariées au sein d’un jeu de données contenant des mesures synthétiques provenant de trois capteurs IoT.

Pour en savoir plus sur le modèle de la forêt d’isolement, consultez l’article original de Liu et al.

Prerequisites

  1. Obtenez un abonnement Microsoft Fabric. Ou, inscrivez-vous pour un essai gratuit de Microsoft Fabric.
  2. Attachez votre carnet à une maison-lac. Sur le côté gauche, sélectionnez Ajouter pour ajouter un lac existant ou créer un lac.
  3. SynapseML est préinstallé dans les runtimes Fabric PySpark (Runtime 1.3 ou ultérieur recommandé). Pour utiliser une version spécifique, voir Installer une version différente de SynapseML sur Fabric.

Importations de bibliothèques

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()

Données d’entrée

# 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

Lire les données

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

Convertissez les colonnes dans les types de données appropriés.

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)

Préparation des données d’apprentissage

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

Préparation des données de test

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

Entraîner le modèle de forêt d'isolement

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)
)

Ensuite, créez un pipeline ML pour entraîner le modèle Isolation Forest.

Pour entraîner le modèle et effectuer des inférences dans le même carnet, l’objet modèle est suffisant. Pour persister et réutiliser le modèle entre les sessions, enregistrez-le auprès de MLflow dans Microsoft Fabric.

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

Effectuer une inférence

Appliquez le modèle entraîné aux données de test :

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

Détecteur d’anomalies prédéfini

Important

Microsoft va retirer le service Azure AI Détecteur d'anomalies le 1er octobre 2026. Depuis le 20 septembre 2023, vous ne pouvez plus créer de nouvelles ressources. Pour une alternative prise en charge, voir Détection des anomalies dans Microsoft Fabric Real-Time Intelligence.

Détecteur d’anomalies Azure AI

  • Statut d’anomalie du dernier point : génère un modèle en utilisant les points précédents et détermine si le dernier point est anormal. Voir le dépôt SynapseML GitHub pour la référence actuelle de l’API.
  • Trouver des anomalies : génère un modèle en utilisant une série entière et trouve des anomalies dans la série. Voir le dépôt SynapseML GitHub pour la référence actuelle de l’API.