このチュートリアルでは、イベントハウスに格納されているサンプル データを使用して、Fabric ノートブックで多変量異常検出モデルをトレーニングする方法について説明します。 その後、KQL クエリセットでトレーニング済みのモデルを使用して、新しいデータをスコア付けし、異常を視覚化します。
背景情報については、「Microsoft Fabric - 概要」の「多変量異常検出」を参照してください。
前提条件
- Microsoft Fabric 対応の 容量を持つ ワークスペース。
- 管理者、共同作成者、またはメンバーワークスペースロール。 環境などの項目を作成するには、このアクセス許可レベルが必要です。
- データベースのあるワークスペース内のイベントハウス。
- サンプル データ ファイル。
- サンプル ノートブック。
パート 1: OneLake の可用性を有効にする
イベントハウスにデータを読み込む前に、 OneLake の可用性 を有効にします。 この設定により、取り込まれたデータが OneLake で使用できるようになり、チュートリアルの後半でノートブックから同じテーブルにアクセスできるようになります。
ワークスペースで、前提条件で作成したイベントハウスを開き、データを格納するデータベースを選択します。
[ データベースの詳細 ] ウィンドウで、 OneLake の可用性 を [オン] に設定します。
パート 2: KQL Python プラグインを有効にする
この手順では、eventhouse で Python プラグインを有効にします。 この手順は、「パート 9: KQL クエリセットの異常を予測する」の KQL クエリセットでPython コードを実行するために必要です。 time-series-anomaly-detector パッケージを含む Python イメージを選択します。
イベントハウスで、リボンの Eventhouse>Plugins を選択します。
[プラグイン] ウィンドウで、言語拡張機能Python[オン] に設定します。
Python 3.11.7 DL を選択します。
完了 を選択します。
パート 3: Spark 環境を作成する
この手順では、多変量異常検出モデルをトレーニングするノートブックを実行する Spark 環境を作成します。 詳細については、「環境の 作成と管理」を参照してください。
ワークスペースから [ + 新しい項目] を選択し、[ 環境] を選択します。
環境名として「
MVAD_ENV」と入力し、[ 作成] を選択します。ライブラリで、[パブリック ライブラリ]を選択します。
[PyPI から追加] を選択します。
検索ボックスに「
time-series-anomaly-detector」と入力します。 [ バージョン ] ボックスに「0.3.9」と入力します。[保存] を選択します。
環境内の [ホーム] タブを選択します。
リボンから、[発行] アイコンを選択します。
すべて公開を選択します。 この手順は、完了するまでに数分かかることがあります。
パート 4: イベントハウスにデータを読み込む
イベントハウスで、データを格納する KQL データベースにカーソルを合わせ、[その他] メニュー [...] を選択します>データを取得します>ローカル ファイル。
[+ 新しいテーブル] を選択し、テーブル名として「
demo_stocks_change」と入力します。アップロード ダイアログで、[ ファイルの参照] を選択し、「 前提条件」でダウンロードしたサンプル データ ファイルをアップロードします。
[次へ] を選択します。
[データの検査]セクションで、[最初の行が列ヘッダー]として[オン]に設定されていることを確認します。
完了 を選択します。
データをアップロードしたら、[閉じる] を選択します。
パート 5: OneLake パスをコピーする
demo_stocks_change テーブルを選択します。 [テーブルの詳細 ペインで、OneLake フォルダー 選択して、OneLake パスをクリップボードにコピーします。 後で使用できるように、テキスト エディターにパスを保存します。
パート 6: ノートブックを準備する
ワークスペースを選択します。
[>>。
[ アップロード] を選択し、[ 前提条件] でダウンロードしたノートブックを選択します。
ノートブックがアップロードされたら、ワークスペースからノートブックを見つけて開くことができます。
上部のリボンで、ワークスペースの 既定 のドロップダウン リストを選択し、前の手順で作成した環境を選択します。
パート 7: ノートブックを実行する
標準パッケージをインポートします。
import numpy as np import pandas as pdSpark では、OneLake ストレージに安全に接続するために ABFSS URI が必要であるため、OneLake URI を ABFSS URI に変換するヘルパー関数を定義します。
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_uriOneLakeTableURIをパート 5 でコピーした OneLake URI に置き換えます。OneLake パスをコピーし、demo_stocks_changeテーブルを 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)次のセルを実行して、トレーニングデータフレームと予測データフレームを準備します。
注
実際の予測は、「 パート 9: KQL クエリセットの異常を予測する」のイベントハウスで実行されます。 本番環境では、通常、新しいストリーミング データをスコアリングします。 このチュートリアルでは、データセットを日付別にトレーニング範囲と予測範囲に分割し、履歴データと受信データをシミュレートします。
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')セルを実行してモデルをトレーニングし、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, )次のセルを実行して、後で KQL Python サンドボックスの予測に使用する登録済みのモデル パスを取得します。
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)最後のセルの出力からモデル URI をコピーします。 パート 9 で使用します。
パート 8: KQL クエリセットを作成する
一般的な情報については、「KQL クエリセットの作成」を参照してください。
- ワークスペースで、[ + 新しい項目>KQL クエリセット] を選択します。
- 「
MultivariateAnomalyDetectionTutorial」と入力し、[ 作成] を選択します。 - OneLake カタログ ウィンドウで、データを格納した KQL データベースを選択します。
- [接続] を選択します。
パート 9: KQL クエリセットの異常を予測する
次の
.create-or-alter functionクエリを実行して、predict_fabric_mvad_fl()ストアド関数を定義します。.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) }次の予測クエリを実行します。
enter your model URI hereを、パート 7: ノートブックの実行の最後にコピーした URI に置き換えます。このクエリでは、トレーニング済みのモデルを使用して 5 つの株式の多変量異常を検出し、結果を
anomalychartとしてレンダリングします。 異常なポイントは最初の株式 (AAPL) に表示されますが、特定の日付の 5 つの株式すべてにおける共同動作の異常を表します。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')
結果の異常グラフは、次の図のようになります。
リソースをクリーンアップする
チュートリアルを完了したら、不要なコストを回避するために作成したリソースを削除します。
- ワークスペースのホームページを参照します。
- このチュートリアルで作成した環境を削除します。
- このチュートリアルで作成したノートブックを削除します。
- このチュートリアルで使用するイベントハウスまたは データベース を削除します。
- このチュートリアルで作成した KQL クエリセットを削除します。