Bemærk
Adgang til denne side kræver godkendelse. Du kan prøve at logge på eller ændre mapper.
Adgang til denne side kræver godkendelse. Du kan prøve at ændre mapper.
Denne artikel viser, hvordan man bruger SynapseML på Apache Spark til detektion af multivariate anomalier. Multivariat anomali-detektion opdager anomalier blandt mange variable eller tidsserier, idet alle indbyrdes korrelationer og afhængigheder mellem de forskellige variable tages i betragtning. I dette scenarie bruger du SynapseML til at træne en isolationsskovmodel til multivariat anomalidetektion, og derefter bruger du den trænede model til at udlede multivariate anomalier i et datasæt med syntetiske målinger fra tre IoT-sensorer.
For at lære mere om isolationsskovmodellen, se den oprindelige artikel af Liu et al.
Forudsætninger
- Få et Microsoft Fabric abonnement. Eller tilmeld dig en gratis Microsoft Fabric prøveperiode.
- Vedhæft din notesbog til et lakehouse. I venstre side skal du vælge Tilføj for at tilføje et eksisterende lakehouse eller oprette et lakehouse.
- SynapseML er forudinstalleret i Fabric PySpark-runtimes (Runtime 1.3 eller senere anbefales). For at bruge en specifik version, se Install a different version of SynapseML on 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()
Inputdata
# 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"
)
)
Cast kolonnerne til de relevante 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)
Forberedelse af 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)
Testdataforberedelse
# filter to data with timestamps within the inference window
df_test = df.filter(
(F.col(timestampColumn) >= inferenceStartTime)
& (F.col(timestampColumn) <= inferenceEndTime)
)
display(df_test)
Train 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)
)
Dernæst skal du oprette en ML-pipeline til at træne Isolation Forest-modellen.
For at træne modellen og udføre slutninger i samme notesbog er modelobjektet tilstrækkeligt. For at bevare og genbruge modellen på tværs af sessioner, registrer den hos MLflow i Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Udfør inferenser
Anvend den trænede model på testdataene:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Forudlavet anomalidetektor
Vigtigt!
Microsoft pensionerer Azure AI Anomaly Detector-tjenesten den 1. oktober 2026. Siden den 20. september 2023 kan du ikke oprette nye ressourcer. For et understøttet alternativ, se Anomaly Detection i Microsoft Fabric Real-Time Intelligence.
- Anomalistatus for det seneste punkt: genererer en model ved at bruge de foregående punkter og afgør, om det seneste punkt er unormalt. Se SynapseML GitHub-repositoriet for aktuel API-reference.
- Find anomalier: genererer en model ved at bruge en hel serie og finder anomalier i serien. Se SynapseML GitHub-repositoriet for aktuel API-reference.