Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Este artigo mostra como usar o SynapseML no Apache Spark para detecção multivariada de anomalias. A detecção de anomalias multivariadas detecta anomalias entre muitas variáveis ou séries temporais, levando em conta todas as inter-correlações e dependências entre as diferentes variáveis. Neste cenário, você usa o SynapseML para treinar um modelo de floresta de isolamento para detecção de anomalias multivariadas e, em seguida, usa 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 de floresta de isolamento, veja o artigo original de Liu et al.
Pré-requisitos
- Obtenha uma assinatura do Microsoft Fabric. Ou inscreva-se para uma avaliação gratuita Microsoft Fabric.
- Anexe seu notebook a um lakehouse. No lado esquerdo, selecione Adicionar para adicionar um lakehouse existente ou criar um.
- O SynapseML vem pré-instalado nos runtimes do Fabric para PySpark (recomenda-se o 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 biblioteca
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"
)
)
Converta 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 treinamento
# 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)
Treinar modelo de floresta de isolamento
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)
)
Em 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, registre-o no MLflow no Microsoft Fabric.
va = VectorAssembler(inputCols=inputCols, outputCol="features")
pipeline = Pipeline(stages=[va, isolationForest])
model = pipeline.fit(df_train)
Executar processo de inferência
Aplique o modelo treinado aos dados do teste:
df_test_pred = model.transform(df_test)
display(df_test_pred)
Detector de Anomalias Premade
Importante
A Microsoft está encerrando o serviço Detector de Anomalias de IA do Azure em 1º de outubro de 2026. Desde 20 de setembro de 2023, você não pode criar novos recursos. Para uma alternativa com suporte, consulte detecção de anomalias no Microsoft Fabric Real-Time Intelligence.
Detector de Anomalias de IA do Azure
- Status de anomalia do ponto mais recente: gera um modelo usando pontos anteriores e determina se o ponto mais recente é anômalo. Veja 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. Veja o repositório SynapseML GitHub para referência atual da API.