Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ten artykuł pokazuje, jak używać SynapseML na Apache Spark do wykrywania wielowymiarowych anomalii. Wykrywanie anomalii wielowymiarowych wykrywa anomalie między wieloma zmiennymi lub szeregami czasowymi, uwzględniając wszystkie wzajemne korelacje i zależności między tymi zmiennymi. W tym scenariuszu używasz SynapseML do trenowania modelu lasu izolacyjnego do wykrywania anomalii wielowariantnych, a następnie używasz wytrenowanego modelu do wnioskowania anomalii wielowymiarowych w zbiorze danych zawierającym syntetyczne pomiary z trzech czujników IoT.
Aby dowiedzieć się więcej o modelu lasu izolacyjnego, zobacz oryginalny artykuł Liu i in.
Wymagania wstępne
- Uzyskaj subskrypcję usługi Microsoft Fabric. Możesz też utworzyć konto bezpłatnej wersji próbnej usługi Microsoft Fabric.
- Dołącz swój notebook do lakehouse. Po lewej stronie wybierz pozycję Dodaj, aby dodać istniejący lakehouse lub utworzyć nowy lakehouse.
- SynapseML jest preinstalowany w środowiskach uruchomieniowych Fabric PySpark (zalecane jest środowisko uruchomieniowe w wersji 1.3 lub nowszej). Aby użyć określonej wersji, zobacz Instalowanie innej wersji SynapseML w usłudze Fabric.
Importowanie biblioteki
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()
Dane wejściowe
# 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
Odczyt danych
df = (
spark.read.format("csv")
.option("header", "true")
.load(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
)
)
Przypisuj kolumny odpowiednim typom danych.
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)
Przygotowywanie danych treningowych
# filter to data with timestamps within the training window
df_train = df.filter(
(F.col(timestampColumn) >= trainingStartTime)
& (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)
Przygotowywanie danych testowych
# filter to data with timestamps within the inference window
df_test = df.filter(
(F.col(timestampColumn) >= inferenceStartTime)
& (F.col(timestampColumn) <= inferenceEndTime)
)
display(df_test)
Trenowanie modelu lasu izolacji
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)
)
Następnie utwórz potok uczenia maszynowego do trenowania modelu Isolation Forest.
Do trenowania modelu i wykonywania wnioskowania w tym samym notesie wystarcza obiekt modelu. Aby utrzymywać i ponownie używać modelu w różnych sesjach, zarejestruj go w MLflow w Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Przeprowadź wnioskowanie
Zastosuj wytrenowany model do danych testowych:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Wstępnie utworzony detektor anomalii
Important
Microsoft wycofuje usługę Detektor anomalii platformy Azure AI 1 października 2026 roku. Od 20 września 2023 roku nie możesz tworzyć nowych zasobów. Wspieraną alternatywę można znaleźć w artykule Wykrywanie anomalii w Microsoft Fabric Real-Time Intelligence.
Narzędzie do wykrywania anomalii w usłudze Azure AI
- Status anomalii najpóźniejszego punktu: generuje model na podstawie poprzednich punktów i określa, czy najnowszy punkt jest anomaliczny. Zobacz repozytorium SynapseML GitHub dla aktualnych źródeł API.
- Znajdź anomalie: generuje model, używając całej serii i znajduje anomalie w tej serii. Zobacz repozytorium SynapseML GitHub dla aktualnych źródeł API.