Detecção multivariada de anomalias com floresta de isolamento

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

  1. Obtenha uma assinatura do Microsoft Fabric. Ou inscreva-se para uma avaliação gratuita Microsoft Fabric.
  2. Anexe seu notebook a um lakehouse. No lado esquerdo, selecione Adicionar para adicionar um lakehouse existente ou criar um.
  3. 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.