Rilevamento di anomalie multivariate

Questa esercitazione illustra come eseguire il training di un modello di rilevamento anomalie multivariato in un notebook di Fabric usando i dati di esempio archiviati in una eventhouse. Si utilizza quindi il modello addestrato in un set di query in KQL per attribuire un punteggio ai nuovi dati e visualizzare le anomalie.

Per informazioni di base, vedi Rilevamento di anomalie multivariate in Microsoft Fabric - panoramica.

Prerequisiti

Parte 1: Attivare la disponibilità di OneLake

Attivare la disponibilità di OneLake prima di caricare i dati nella eventhouse. Questa impostazione rende disponibili i dati inseriti in OneLake in modo che sia possibile accedere alla stessa tabella da un notebook più avanti nell'esercitazione.

  1. Nell'area di lavoro aprire la eventhouse creata nei prerequisiti e quindi selezionare il database in cui archiviare i dati.

  2. Nel riquadro Dettagli database impostare La disponibilità di OneLakesu Sì.

    Schermata dell'attivazione della disponibilità di OneLake nell'eventhouse.

Parte 2: Attivare il plug-in KQL Python

In questo passaggio si attiva il plug-in Python nell'eventhouse. Questo passaggio è necessario per eseguire il codice Python nel set di query KQL nella parte 9: Stimare le anomalie in un set di query KQL. Selezionare l'immagine Python che include il pacchetto time series-anomaly-detector.

  1. Nella eventhouse, selezionare Eventhouse>Plugins sulla barra multifunzione.

  2. Nel riquadro Plugin, imposta estensione del linguaggio Python su Attivata.

  3. Selezionare Python DL 3.11.7.

  4. Selezionare Fatto.

    Screenshot che mostra come abilitare il pacchetto DL Python 3.11.7 nella eventhouse.

Parte 3: Creare un ambiente Spark

In questo passaggio si crea un ambiente Spark per eseguire il notebook che addestra il modello multivariato di rilevamento delle anomalie. Per altre informazioni, vedere Creare e gestire ambienti.

  1. Nell'area di lavoro selezionare + Nuovo elemento e quindi selezionare Ambiente.

    Screenshot del riquadro Ambiente nella finestra Nuovo elemento.

  2. Immettere MVAD_ENV per il nome dell'ambiente e quindi selezionare Crea.

  3. In Cataloghi selezionare Cataloghi pubblici.

  4. Selezionare Aggiungi da PyPi.

  5. Nella casella di ricerca immettere time-series-anomaly-detector. Nella casella Versione immettere 0.3.9.

  6. Seleziona Salva.

    Screenshot dell'aggiunta del pacchetto PyPI all'ambiente Spark.

  7. Selezionare la scheda Home nell’ambiente.

  8. Seleziona l’icona Pubblica nella barra multifunzione.

  9. Selezionare Pubblica tutto. Questo passaggio può richiedere alcuni minuti.

    Screenshot della pubblicazione dell'ambiente.

Parte 4: Caricare i dati nella eventhouse

  1. Nella casa eventi passare il puntatore del mouse sul database KQL in cui archiviare i dati e quindi selezionare Altro menu [...]>Recuperare i dati>File locale.

    Screenshot del recupero dei dati dal file locale.

  2. Selezionare + Nuova tabella e immettere demo_stocks_change come nome della tabella.

  3. Nella finestra di dialogo di caricamento selezionare Cerca file e caricare il file di dati di esempio scaricato in Prerequisiti.

  4. Selezionare Avanti.

  5. Nella sezione Controlla i dati verificare che l'intestazione della prima riga sia impostata su .

  6. Selezionare Fine.

  7. Al termine del caricamento selezionare Chiudi.

Parte 5: Copiare il percorso di OneLake

Selezionare la demo_stocks_change tabella. Nel riquadro Dettagli tabella selezionare cartella OneLake per copiare il percorso di OneLake negli Appunti. Salvare il percorso in un editor di testo per usarlo in un secondo momento.

Screenshot che mostra la copia del percorso di OneLake.

Parte 6: Preparare il notebook

  1. Selezionare l'area di lavoro.

  2. Selezionare Importa>notebook>da questo computer.

  3. Selezionare Carica e scegliere il notebook scaricato in Prerequisiti.

  4. Dopo aver caricato il notebook, è possibile trovare e aprire il notebook dall'area di lavoro.

  5. Nella barra multifunzione superiore selezionare l'elenco a discesa Area di lavoro predefinito e quindi selezionare l'ambiente creato nel passaggio precedente.

    Screenshot della selezione dell'ambiente nel notebook.

Parte 7: Esegui il notebook

  1. Importare pacchetti standard.

    import numpy as np
    import pandas as pd
    
  2. Spark richiede un URI ABFSS per connettersi in modo sicuro all'archiviazione OneLake, quindi definire una funzione helper che converte l'URI di OneLake in un URI ABFSS.

    def convert_onelake_to_abfss(onelake_uri):
        if not onelake_uri.startswith('https://'):
            raise ValueError("Invalid OneLake URI. It should start with 'https://'.")
        uri_without_scheme = onelake_uri[8:]
        parts = uri_without_scheme.split('/')
        if len(parts) < 3:
            raise ValueError("Invalid OneLake URI format.")
        container_name = parts[1]
        path = '/'.join(parts[2:])
        abfss_uri = f"abfss://{container_name}@{parts[0]}/{path}"
        return abfss_uri
    
  3. Sostituire OneLakeTableURI con l'URI di OneLake copiato nella parte 5: Copiare il percorso di OneLake e quindi caricare la demo_stocks_change tabella in un dataframe pandas.

    onelake_uri = "OneLakeTableURI"  # Replace with your OneLake table URI.
    abfss_uri = convert_onelake_to_abfss(onelake_uri)
    print(abfss_uri)
    
    df = spark.read.format('delta').load(abfss_uri)
    df = df.toPandas()
    df['Date'] = pd.to_datetime(df['Date'])
    df = df.set_index('Date').sort_index()
    print(df.shape)
    df.head(3)
    
  4. Eseguire le celle seguenti per preparare i dataframe di training e predizione.

    Nota

    Le stime effettive vengono eseguite nella eventhouse nella parte 9: Stimare le anomalie in un set di query KQL. In uno scenario di produzione, in genere si assegna un punteggio ai nuovi dati in streaming. In questo tutorial, il set di dati viene suddiviso per data in intervalli di training e previsione per simulare dati storici e dati in arrivo.

    features_cols = ['AAPL', 'AMZN', 'GOOG', 'MSFT', 'SPY']
    cutoff_date = pd.Timestamp('2023-01-01')
    
    train_df = df.loc[df.index < cutoff_date, features_cols]
    print(train_df.shape)
    train_df.head(3)
    
    train_len = len(train_df)
    predict_len = len(df) - train_len
    print(f'Total samples: {len(df)}. Split to {train_len} for training, {predict_len} for testing')
    
  5. Eseguite le celle per addestrare il modello e salvarlo nel registro dei modelli MLflow di Fabric.

    from anomaly_detector import MultivariateAnomalyDetector
    model = MultivariateAnomalyDetector()
    
    sliding_window = 200
    params = {"sliding_window": sliding_window}
    
    model.fit(train_df, params=params)
    
    model_name = "mvad_5_stocks_model"
    
    import mlflow
    
    with mlflow.start_run():
        mlflow.log_params(params)
        mlflow.set_tag("Training Info", "MVAD on 5 Stocks Dataset")
    
        model_info = mlflow.pyfunc.log_model(
            python_model=model,
            artifact_path="mvad_artifacts",
            registered_model_name=model_name,
        )
    
  6. Esegui la cella seguente per ottenere il percorso del modello registrato da usare successivamente per la previsione nella sandbox Python KQL.

    from mlflow.tracking import MlflowClient
    
    client = MlflowClient()
    mvs = client.search_model_versions(f"name='{model_name}'")
    latest = max(mvs, key=lambda v: v.creation_timestamp)
    model_abfss = latest.source
    print(model_abfss)
    
  7. Copiare l'URI del modello dall'output dell'ultima cella. La si usa nella parte 9.

Parte 8: Creare un insieme di query KQL

Per informazioni generali, vedere Creare un set di query KQL.

  1. Nell'area di lavoro selezionare + Nuovo elemento>KQL Queryset.
  2. Immettere MultivariateAnomalyDetectionTutoriale quindi selezionare Crea.
  3. Nella finestra del catalogo OneLake selezionare il database KQL in cui sono stati archiviati i dati.
  4. Selezionare Connetti.

Parte 9: Prevedere le anomalie in un set di query KQL

  1. Esegui la seguente query .create-or-alter function per definire la predict_fabric_mvad_fl() funzione memorizzata:

    .create-or-alter function with (folder = "Packages\\ML", docstring = "Predict MVAD model in Microsoft Fabric")
    predict_fabric_mvad_fl(samples:(*), features_cols:dynamic, artifacts_uri:string, trim_result:bool=false)
    {
        let s = artifacts_uri;
        let artifacts = bag_pack('MLmodel', strcat(s, '/MLmodel;impersonate'), 'conda.yaml', strcat(s, '/conda.yaml;impersonate'),
                                 'requirements.txt', strcat(s, '/requirements.txt;impersonate'), 'python_env.yaml', strcat(s, '/python_env.yaml;impersonate'),
                                 'python_model.pkl', strcat(s, '/python_model.pkl;impersonate'));
        let kwargs = bag_pack('features_cols', features_cols, 'trim_result', trim_result);
        let code = ```if 1:
            import os
            import shutil
            import mlflow
            work_dir = os.environ.get("UPLOAD_PATH")
            model_dir = work_dir + '/mvad_model'
            model_data_dir = model_dir + '/data'
            os.mkdir(model_dir)
            shutil.move(work_dir + '/MLmodel', model_dir)
            shutil.move(work_dir + '/conda.yaml', model_dir)
            shutil.move(work_dir + '/requirements.txt', model_dir)
            shutil.move(work_dir + '/python_env.yaml', model_dir)
            shutil.move(work_dir + '/python_model.pkl', model_dir)
            features_cols = kargs["features_cols"]
            trim_result = kargs["trim_result"]
            test_data = df[features_cols]
            model = mlflow.pyfunc.load_model(model_dir)
            predictions = model.predict(test_data)
            predict_result = pd.DataFrame(predictions)
            samples_offset = len(df) - len(predict_result)        # this model doesn't output predictions for the first sliding_window-1 samples
            if trim_result:                                       # trim the prefix samples
                result = df[samples_offset:]
                result.iloc[:,-4:] = predict_result.iloc[:, 1:]   # no need to copy 1st column which is the timestamp index
            else:
                result = df                                       # output all samples
                result.iloc[samples_offset:,-4:] = predict_result.iloc[:, 1:]
            ```;
        samples
        | evaluate python(typeof(*), code, kwargs, external_artifacts=artifacts)
    }
    
  2. Esegui la seguente query di previsione. Sostituisci enter your model URI here con l'URI che hai copiato alla fine di Parte 7: Esegui il notebook.

    La query rileva anomalie multivariate nei cinque titoli utilizzando il modello addestrato e quindi visualizza i risultati come anomalychart. I punti anomali vengono visualizzati sul primo titolo (AAPL), ma rappresentano anomalie nel comportamento congiunto di tutti e cinque i titoli in una data specifica.

    let cutoff_date=datetime(2023-01-01);
    let num_predictions=toscalar(demo_stocks_change | where Date >= cutoff_date | count);   //  number of latest points to predict
    let sliding_window=200;                                                                 //  should match the window that was set for model training
    let prefix_score_len = sliding_window/2+min_of(sliding_window/2, 200)-1;
    let num_samples = prefix_score_len + num_predictions;
    demo_stocks_change
    | top num_samples by Date desc
    | order by Date asc
    | extend is_anomaly=bool(false), score=real(null), severity=real(null), interpretation=dynamic(null)
    | invoke predict_fabric_mvad_fl(pack_array('AAPL', 'AMZN', 'GOOG', 'MSFT', 'SPY'),
                // NOTE: Update artifacts_uri to model path
                artifacts_uri='enter your model URI here',
                trim_result=true)
    | summarize Date=make_list(Date), AAPL=make_list(AAPL), AMZN=make_list(AMZN), GOOG=make_list(GOOG), MSFT=make_list(MSFT), SPY=make_list(SPY), anomaly=make_list(toint(is_anomaly))
    | render anomalychart with(anomalycolumns=anomaly, title='Stock price changes in % with anomalies')
    

Il grafico delle anomalie risultante è simile all'immagine seguente:

Screenshot dell'output delle anomalie multivariate.

Pulire le risorse

Al termine dell'esercitazione, eliminare le risorse create per evitare costi non necessari:

  1. Naviga alla home page dell'area di lavoro.
  2. Eliminare l’ambiente creato in questa esercitazione.
  3. Eliminare il notebook creato in questo tutorial.
  4. Eliminare la eventhouse o il database usato in questa esercitazione.
  5. Eliminare il set di query KQL creato in questa esercitazione.