Wielowymiarowe wykrywanie anomalii

W tym samouczku pokazano, jak trenować model wykrywania anomalii wielowymiarowych w notesie w usłudze Fabric przy użyciu przykładowych danych przechowywanych w Eventhouse. Następnie użyj wytrenowanego modelu w zestawie zapytań KQL, aby ocenić nowe dane i zwizualizować anomalie.

Aby uzyskać informacje wprowadzające, zobacz Wykrywanie anomalii wielowariancyjnych w Microsoft Fabric — omówienie.

Wymagania wstępne

Część 1: Włącz dostępność OneLake

Włącz dostępność usługi OneLake przed załadowaniem danych do magazynu zdarzeń. To ustawienie udostępnia pozyskane dane w usłudze OneLake, dzięki czemu można uzyskać dostęp do tej samej tabeli z notesu w dalszej części tego samouczka.

  1. W obszarze roboczym otwórz magazyn zdarzeń utworzony w wymaganiach wstępnych, a następnie wybierz bazę danych, w której chcesz przechowywać dane.

  2. W panelu Szczegóły bazy danych ustaw parametr Dostępność usługi OneLake na Włączone.

    Zrzut ekranu przedstawiający włączanie dostępności usługi OneLake w Twoim Eventhouse.

Część 2. Włączanie wtyczki KQL Python

W tym kroku włączasz wtyczkę języka Python w swoim Eventhouse. Ten krok jest wymagany do uruchomienia kodu Python w zestawie zapytań KQL w części 9: Przewidywanie anomalii w zestawie zapytań KQL. Wybierz obraz Python zawierający pakiet narzędzia do wykrywania anomalii szeregów czasowych.

  1. W środowisku Eventhouse> na wstążce wybierz pozycję Wtyczki.

  2. W panelu Wtyczki ustaw rozszerzenie języka Python na Włączone.

  3. Wybierz Python 3.11.7 DL.

  4. Wybierz pozycję Gotowe.

    Zrzut ekranu pokazujący, jak włączyć pakiet Python 3.11.7 DL w Eventhouse.

Część 3. Tworzenie środowiska Spark

W tym kroku utworzysz środowisko platformy Spark, aby uruchomić notes, który trenuje wielowariantowy model wykrywania anomalii. Aby uzyskać więcej informacji, zobacz Tworzenie środowisk i zarządzanie nimi.

  1. W obszarze roboczym wybierz pozycję + Nowy element, a następnie wybierz pozycję Środowisko.

    Zrzut ekranu kafelka

  2. Wprowadź MVAD_ENV nazwę środowiska, a następnie wybierz pozycję Utwórz.

  3. W obszarze Biblioteki wybierz pozycję Biblioteki publiczne.

  4. Wybierz Dodaj z PyPI.

  5. W polu wyszukiwania wpisz time-series-anomaly-detector. W polu Wersja wprowadź wartość 0.3.9.

  6. Wybierz pozycję Zapisz.

    Zrzut ekranu przedstawiający dodawanie pakietu PyPI do środowiska Spark.

  7. Wybierz kartę Narzędzia główne w środowisku.

  8. Wybierz ikonę Publikuj na wstążce.

  9. Wybierz opcję Publikuj wszystko. Wykonanie tego kroku może potrwać kilka minut.

    Zrzut ekranu przedstawiający publikowanie środowiska.

Część 4. Ładowanie danych do magazynu zdarzeń

  1. W usłudze Eventhouse najedź kursorem na bazę danych KQL, w której chcesz przechowywać dane, a następnie wybierz menu Więcej [...]>Pobierz dane>Plik lokalny.

    Zrzut ekranu przedstawiający pobieranie danych z pliku lokalnego.

  2. Wybierz pozycję + Nowa tabela i wprowadź jako demo_stocks_change nazwę tabeli.

  3. W oknie dialogowym przekazywania plików wybierz pozycję Przeglądaj pliki i prześlij przykładowy plik danych pobrany w sekcji Wymagania wstępne.

  4. Wybierz Dalej.

  5. W sekcji Inspekcja danych sprawdź, czy Pierwszy wiersz to nagłówek kolumny jest ustawiony na Włącz.

  6. Wybierz Zakończ.

  7. Po przekazaniu danych wybierz pozycję Zamknij.

Część 5. Kopiowanie ścieżki OneLake

Wybierz tabelę demo_stocks_change . W okienku Szczegóły tabeli wybierz folder OneLake, aby skopiować ścieżkę OneLake do schowka. Zapisz ścieżkę w edytorze tekstów do późniejszego użycia.

Zrzut ekranu przedstawiający kopiowanie ścieżki OneLake.

Część 6: Przygotuj notebook

  1. Wybierz obszar roboczy.

  2. Wybierz pozycję Importuj>Notatnik>z tego komputera.

  3. Wybierz pozycję Przekaż i wybierz notes pobrany w sekcji Wymagania wstępne.

  4. Po przesłaniu notatnika możesz znaleźć i otworzyć go ze swojego obszaru roboczego.

  5. Na wstążce u góry wybierz listę rozwijaną Domyślny obszar roboczy, a następnie wybierz środowisko utworzone w poprzednim kroku.

    Zrzut ekranu przedstawiający wybieranie środowiska w notesie.

Część 7: Uruchom notebook

  1. Importuj standardowe pakiety.

    import numpy as np
    import pandas as pd
    
  2. Platforma Spark wymaga identyfikatora URI ABFSS, aby bezpiecznie nawiązać połączenie z magazynem OneLake, dlatego zdefiniuj funkcję pomocnika, która konwertuje identyfikator URI usługi OneLake na identyfikator 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. Zastąp element OneLakeTableURI identyfikatorem URI OneLake skopiowanym w części 5: Skopiuj ścieżkę OneLake, a następnie załaduj demo_stocks_change tabelę do ramki danych biblioteki 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. Uruchom następujące komórki, aby przygotować ramki danych trenowania i przewidywania.

    Uwaga

    Rzeczywiste przewidywania są wykonywane w Eventhouse w Część 9: Przewidywanie anomalii w zbiorze zapytań KQL. W scenariuszu produkcyjnym zazwyczaj oceniasz nowe dane przesyłane strumieniowo. W tym samouczku zestaw danych jest podzielony według daty na zakresy trenowania i przewidywania w celu symulowania danych historycznych i przychodzących.

    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. Uruchom komórki, aby wytrenować model i zapisać go w rejestrze modeli Fabric MLflow.

    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. Uruchom następującą komórkę, aby uzyskać zarejestrowaną ścieżkę modelu używaną później do przewidywania w piaskownicy 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. Skopiuj identyfikator URI modelu z danych wyjściowych ostatniej komórki. Używasz go w części 9.

Część 8: Tworzenie zestawu zapytań KQL

Aby uzyskać ogólne informacje, zobacz Tworzenie zestawu zapytań KQL.

  1. W obszarze roboczym wybierz pozycję + Nowy element>Zestaw zapytań KQL.
  2. Wprowadź ciąg MultivariateAnomalyDetectionTutorial, a następnie wybierz pozycję Utwórz.
  3. W oknie katalogu OneLake wybierz bazę danych KQL, w której zapisano dane.
  4. Wybierz pozycję Połącz.

Część 9. Przewidywanie anomalii w zestawie zapytań KQL

  1. Uruchom następujące .create-or-alter function zapytanie, aby zdefiniować przechowywaną predict_fabric_mvad_fl() funkcję:

    .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. Uruchom następujące zapytanie przewidywania. Zastąp enter your model URI here adresem URI skopiowanym pod koniec Część 7: Uruchamianie notesu.

    Zapytanie wykrywa anomalie wielowariancji w pięciu akcjach przy użyciu wytrenowanego modelu, a następnie renderuje wyniki jako anomalychart. Punkty odstające są wyświetlane przy pierwszej akcji (AAPL), ale reprezentują anomalie w łącznym zachowaniu wszystkich pięciu akcji w danym dniu.

    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')
    

Wynikowy wykres anomalii przypomina następujący obraz:

Zrzut ekranu wyników anomalii wielowariancyjnych.

Czyszczenie zasobów

Po ukończeniu samouczka usuń utworzone zasoby, aby uniknąć niepotrzebnych kosztów:

  1. Wejdź na stronę główną obszaru roboczego.
  2. Usuń środowisko utworzone w tym samouczku.
  3. Usuń notatnik utworzony w tym samouczku.
  4. Usuń magazyn zdarzeń lub bazę danych używaną w tym samouczku.
  5. Usuń zestaw zapytań KQL utworzony w tym samouczku.