Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
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
- Haal een Microsoft Fabric-abonnement op. Of meld u aan voor een gratis proefversie van Microsoft Fabric.
- Koppel uw notitieblok aan een lakehouse. Selecteer aan de linkerkant toevoegen om een bestaand lakehouse toe te voegen of een lakehouse te maken.
- 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.
- 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.