Merk
Tilgang til denne siden krever autorisasjon. Du kan prøve å logge på eller endre kataloger.
Tilgang til denne siden krever autorisasjon. Du kan prøve å endre kataloger.
Denne artikkelen viser hvordan man bruker SynapseML på Apache Spark for deteksjon av multivariate anomalier. Multivariat anomalideteksjon oppdager anomalier mellom mange variabler eller tidsserier, og tar hensyn til alle interkorrelasjoner og avhengigheter mellom de ulike variablene. I dette scenariet bruker du SynapseML til å trene en isolasjonsskogmodell for deteksjon av multivariate anomalier, og deretter bruker du den trente modellen til å utlede multivariate anomalier i et datasett som inneholder syntetiske målinger fra tre IoT-sensorer.
For å lære mer om isolasjonsskogmodellen, se den opprinnelige artikkelen av Liu et al.
Forutsetninger
- Skaff deg et abonnement Microsoft Fabric. Eller meld deg på en gratis prøveperiode Microsoft Fabric.
- Legg notatblokken til et lakehouse. På venstre side velger du Legg til for å legge til et eksisterende innsjøhus eller opprette et innsjøhus.
- SynapseML er forhåndsinstallert i Fabric PySpark-kjøretider (Runtime 1.3 eller nyere anbefales). For å bruke en spesifikk versjon, se Installer en annen versjon av SynapseML på Fabric.
Bibliotekimport
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()
Inndata
# 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
Les data
df = (
spark.read.format("csv")
.option("header", "true")
.load(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
)
)
Kast kolonnene til de riktige datatypene.
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 av treningsdata
# filter to data with timestamps within the training window
df_train = df.filter(
(F.col(timestampColumn) >= trainingStartTime)
& (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)
Forberedelse av testdata
# 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)
)
Deretter lager du en ML-pipeline for å trene Isolation Forest-modellen.
For å trene modellen og utføre slutninger i samme notatbok, er modellobjektet tilstrekkelig. For å opprettholde og gjenbruke modellen på tvers av økter, registrer den hos MLflow i Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Utfør slutninger
Bruk den trente modellen på testdataene:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Forhåndslaget anomalidetektor
Viktig!
Microsoft legger ned tjenesten Azure AI Anomaly Detector 1. oktober 2026. Siden 20. september 2023 kan du ikke opprette nye ressurser. For et støttet alternativ, se Anomalideteksjon i Microsoft Fabric Real-Time Intelligence.
- Anomalistatus for siste punkt: genererer en modell ved å bruke foregående punkter og avgjør om det siste punktet er anomalie. Se SynapseML GitHub-repositoriet for oppdatert API-referanse.
- Finn anomalier: genererer en modell ved å bruke en hel serie og finner anomalier i serien. Se SynapseML GitHub-repositoriet for oppdatert API-referanse.