Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Questo articolo mostra come utilizzare SynapseML su Apache Spark per la rilevazione di anomalie multivariate. Il rilevamento di anomalie multivariate rileva anomalie tra molte variabili o serie temporali, tenendo conto di tutte le inter-correlazioni e dipendenze tra le diverse variabili. In questo scenario, si utilizza SynapseML per addestrare un modello di foresta isolata per il rilevamento di anomalie multivariate, e poi si utilizza il modello addestrato per dedurre anomalie multivariate all'interno di un dataset contenente misurazioni sintetiche da tre sensori IoT.
Per saperne di più sul modello della foresta di isolamento, consulta l'articolo originale di Liu et al.
Prerequisiti
- Ottieni un abbonamento a Microsoft Fabric. Oppure, registrati per una versione di prova gratuita di Microsoft Fabric.
- Collega il notebook a un lakehouse. Sul lato sinistro, selezionare Aggiungi per aggiungere un lakehouse esistente o creare un lakehouse.
- SynapseML è preinstallato nei runtime di Fabric PySpark (consigliato il runtime 1.3 o successivo). Per utilizzare una versione specifica, vedi Installa una versione diversa di SynapseML su Fabric.
Importazioni di raccolte
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()
Dati di input
# 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
Leggere i dati
df = (
spark.read.format("csv")
.option("header", "true")
.load(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
)
)
Converti le colonne nei tipi di dati appropriati.
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)
Preparazione dei dati di addestramento
# filter to data with timestamps within the training window
df_train = df.filter(
(F.col(timestampColumn) >= trainingStartTime)
& (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)
Preparazione dei dati di 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)
Addestrare il modello Isolation Forest
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)
)
Successivamente, crea una pipeline di ML per addestrare il modello Isolation Forest.
Per addestrare il modello ed eseguire inferenze nello stesso quaderno, l'oggetto modello è sufficiente. Per persistere e riutilizzare il modello tra le sessioni, registralo con MLflow in Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Eseguire l'inferenza
Applica il modello addestrato ai dati di test:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Rilevatore di anomalie preconfezionato
Importante
Microsoft ritirerà il servizio Rilevamento anomalie di Azure AI il 1° ottobre 2026. Dal 20 settembre 2023, non puoi creare nuove risorse. Per un'alternativa supportata, vedere Rilevamento delle anomalie in Microsoft Fabric Real-Time Intelligence.
Rilevamento anomalie di Azure AI
- Stato anomalico dell'ultimo punto: genera un modello utilizzando i punti precedenti e determina se l'ultimo punto è anomalo. Consulta il repository SynapseML GitHub per il riferimento API attuale.
- Trova anomalie: genera un modello utilizzando un'intera serie e trova anomalie nella serie. Consulta il repository SynapseML GitHub per il riferimento API attuale.
Contenuto correlato
- Come creare un motore di ricerca con SynapseML
- Come usare Gli strumenti SynapseML e Foundry per il rilevamento delle anomalie multivariate - Analizzare le serie temporali
Come usare SynapseML per ottimizzare gli iperparametri