Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Den här artikeln visar hur man använder SynapseML på Apache Spark för multivariat avvikelsedetektion. Multivariat avvikelsedetektion upptäcker avvikelser bland många variabler eller tidsserier, med hänsyn till alla inbördes korrelationer och beroenden mellan de olika variablerna. I detta scenario använder du SynapseML för att träna en isolationsskogsmodell för multivariat anomalidetektion, och sedan använder du den tränade modellen för att härleda multivariata anomalier inom en datamängd som innehåller syntetiska mätningar från tre IoT-sensorer.
För att lära dig mer om isolationsskogsmodellen, se den ursprungliga artikeln av Liu et al.
Förutsättningar
- Skaffa en Microsoft Fabric-prenumeration. Eller registrera dig för en kostnadsfri utvärderingsversion av Microsoft Fabric.
- Anslut din anteckningsbok till ett lakehouse. Till vänster väljer du Lägg till för att lägga till ett befintligt sjöhus eller skapa ett sjöhus.
- SynapseML är förinstallerat i Fabric PySpark-runtimes (Runtime 1.3 eller senare rekommenderas). För att använda en specifik version, se Installera en annan version av SynapseML på Fabric.
Biblioteksimport
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()
Inmatningsdata
# 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
Läs data
df = (
spark.read.format("csv")
.option("header", "true")
.load(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
)
)
Konvertera kolumnerna till lämpliga datatyper.
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)
Förberedelse av träningsdata
# filter to data with timestamps within the training window
df_train = df.filter(
(F.col(timestampColumn) >= trainingStartTime)
& (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)
Testdataförberedelse
# filter to data with timestamps within the inference window
df_test = df.filter(
(F.col(timestampColumn) >= inferenceStartTime)
& (F.col(timestampColumn) <= inferenceEndTime)
)
display(df_test)
Träna Isolation Forest-modellen
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)
)
Därefter, skapa en ML-pipeline för att träna Isolation Forest-modellen.
För att träna modellen och utföra inferenser i samma anteckningsbok är modellobjektet tillräckligt. För att behålla och återanvända modellen över sessioner, registrera den i MLflow i Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Utföra slutsatsdragning
Applicera den tränade modellen på testdatan:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Förbyggd anomalidetektor
Important
Microsoft lägger ner Azure AI-avvikelseidentifiering-tjänsten den 1 oktober 2026. Sedan den 20 september 2023 kan du inte skapa nya resurser. För ett alternativ som stöds, se Anomaly detection in Microsoft Fabric Real-Time Intelligence.
Azure AI-avvikelseidentifiering
- Avvikelsestatus för senaste punkt: genererar en modell genom att använda föregående punkter och avgör om den senaste punkten är avvikande. Se SynapseML GitHub-arkivet för aktuell API-referens.
- Hitta anomalier: genererar en modell genom att använda en hel serie och hittar anomalier i serien. Se SynapseML GitHub-arkivet för aktuell API-referens.