Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Este artigo mostra como usar o SynapseML no Apache Spark para deteção de anomalias multivariadas. A deteção de anomalias multivariadas deteta anomalias entre muitas variáveis ou séries temporais, tendo em conta todas as inter-correlações e dependências entre as diferentes variáveis. Neste cenário, utiliza-se o SynapseML para treinar um modelo de floresta de isolamento para deteção de anomalias multivariadas, e depois utiliza-se o modelo treinado para inferir anomalias multivariadas dentro de um conjunto de dados contendo medições sintéticas de três sensores IoT.
Para saber mais sobre o modelo da floresta de isolamento, consulte o artigo original de Liu et al.
Pré-requisitos
- Obtenha uma assinatura do Microsoft Fabric. Ou inscreva-se para uma avaliação gratuita do Microsoft Fabric.
- Ligue o seu bloco de notas a uma casa no lago. No lado esquerdo, selecione Adicionar para adicionar uma casa de lago existente ou criar uma casa de lago.
- O SynapseML está pré-instalado em runtimes do Fabric PySpark (recomendado Runtime 1.3 ou posterior). Para usar uma versão específica, veja Instalar uma versão diferente do SynapseML no Fabric.
Importações de bibliotecas
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()
Dados de entrada
# 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
Ler dados
df = (
spark.read.format("csv")
.option("header", "true")
.load(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/generated_sample_mvad_data.csv"
)
)
Conjure as colunas para os tipos de dados apropriados.
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)
Preparação de dados de treino
# filter to data with timestamps within the training window
df_train = df.filter(
(F.col(timestampColumn) >= trainingStartTime)
& (F.col(timestampColumn) <= trainingEndTime)
)
display(df_train)
Preparação de dados de teste
# filter to data with timestamps within the inference window
df_test = df.filter(
(F.col(timestampColumn) >= inferenceStartTime)
& (F.col(timestampColumn) <= inferenceEndTime)
)
display(df_test)
Modelo da Floresta de Isolamento de Comboios
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)
)
De seguida, crie um pipeline de ML para treinar o modelo Isolation Forest.
Para treinar o modelo e realizar inferências no mesmo caderno, o objeto modelo é suficiente. Para persistir e reutilizar o modelo entre sessões, registe-o no MLflow no Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Realizar inferência
Aplique o modelo treinado aos dados de teste:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Detetor de Anomalias Pré-Fabricado
Importante
A Microsoft vai retirar o serviço Detector de anomalias de IA do Azure a 1 de outubro de 2026. Desde 20 de setembro de 2023, não é possível criar novos recursos. Para uma alternativa suportada, consulte Deteção de anomalias no Microsoft Fabric Real-Time Intelligence.
- Estado de anomalia do ponto mais recente: gera um modelo usando pontos anteriores e determina se o ponto mais recente é anómalo. Consulte o repositório SynapseML GitHub para referência atual da API.
- Encontrar anomalias: gera um modelo usando uma série inteira e encontra anomalias na série. Consulte o repositório SynapseML GitHub para referência atual da API.