Multivariate anomaliedetectie met isolatiebos

Dit artikel laat zien hoe je SynapseML gebruikt op Apache Spark voor multivariate anomaliedetectie. Multivariate anomaliedetectie detecteert anomalieën tussen vele variabelen of tijdreeksen, waarbij alle onderlinge correlaties en afhankelijkheden tussen de verschillende variabelen worden meegenomen. In dit scenario gebruik je SynapseML om een isolatieforestmodel te trainen voor multivariate anomaliedetectie, en vervolgens gebruik je het getrainde model om multivariate anomalieën af te leiden binnen een dataset met synthetische metingen van drie IoT-sensoren.

Voor meer informatie over het isolatiebosmodel, zie het originele artikel van Liu et al.

Vereiste voorwaarden

  1. Haal een Microsoft Fabric-abonnement op. Of meld u aan voor een gratis proefversie van Microsoft Fabric.
  2. Koppel uw notitieblok aan een lakehouse. Selecteer aan de linkerkant toevoegen om een bestaand lakehouse toe te voegen of een lakehouse te maken.
  3. SynapseML is vooraf geïnstalleerd in Fabric PySpark-runtimes (Runtime 1.3 of later wordt aanbevolen). Om een specifieke versie te gebruiken, raadpleeg Een andere versie van SynapseML installeren op Fabric.

Bibliotheekimport

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

Invoergegevens

# 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

Gegevens lezen

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

Verdeel de kolommen naar de juiste datatypes.

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)

Trainingsgegevensvoorbereiding

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

Gegevensvoorbereiding testen

# 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-model trainen

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

Maak vervolgens een ML-pijplijn om het Isolation Forest-model te trainen.

Voor het trainen van het model en het uitvoeren van inferencing in hetzelfde notebook is het modelobject voldoende. Om het model te behouden en opnieuw te gebruiken over sessies heen, registreer je het bij MLflow in Microsoft Fabric.

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

Inferentie uitvoeren

Pas het getrainde model toe op de testgegevens:

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

Voorgeconfigureerde anomaliedetector

Belangrijk

Microsoft stopt de Azure AI Anomaly Detector-service op 1 oktober 2026. Sinds 20 september 2023 kun je geen nieuwe bronnen meer aanmaken. Voor een ondersteund alternatief, zie Anomaly Detection in Microsoft Fabric Real-Time Intelligence.

Azure AI Anomaly Detector

  • Anomaliestatus van het laatste punt: genereert een model door gebruik te maken van voorgaande punten en bepaalt of het laatste punt anomaal is. Zie de SynapseML GitHub-repository voor de actuele API-referentie.
  • Vind anomalieën: genereert een model door een hele serie te gebruiken en vindt anomalieën in de serie. Zie de SynapseML GitHub-repository voor de actuele API-referentie.