メインコンテンツまでスキップ

Feature Views を使用してモデルをトレーニングする

備考

プレビュー

この機能は パブリック プレビュー段階です。ワークスペース管理者は、 プレビュー ページからこの機能へのアクセスを制御できます。Databricksのプレビューを管理するを参照してください。

特徴量ビューを使用すると、ポイントインタイムに正しい特徴量計算と、推論時の自動特徴量ルックアップでモデルをトレーニングすることができます。Feature Views の定義に関する情報については、「Feature Views」を参照してください。

要件​

APIメソッド​

create_training_set()​

特徴量ビューを作成したら、次のステップはモデルのトレーニングデータを作成することです。これを行うには、ラベル付けされたデータセットをcreate_training_setに渡します。これにより、各特徴量の値が自動的にポイントインタイムで正確に計算されます。

例えば:

Python
FeatureEngineeringClient.create_training_set(
df: DataFrame, # DataFrame with training data
features: Optional[List[Feature]], # List of Feature objects
label: Union[str, List[str], None], # Label column name(s)
exclude_columns: Optional[List[str]] = None, # Optional: columns to exclude
) -> TrainingSet

TrainingSet.load_dfを呼び出して、元のトレーニング データをポイントインタイムの動的コンピュート機能と結合します。

引数dfは、以下の要件を満たす必要があります。

  • 機能定義によって参照されるすべてのエンティティ列を含める必要があります。
  • フィーチャ定義で参照される時系列列が含まれている必要があります。
  • いずれかの RequestSource スキーマで宣言されたすべての列を含める必要があります。タイプは宣言されたスキーマに対して検証されます。不一致の場合、エラーを発生させます(暗黙的なキャストなし)。
  • ラベル列が含まれている必要があります。
  • エンティティ列名、時系列列名、およびリクエスト機能列名のセットは、すべてのソースにおいてグローバルに一意である必要があります。

ポイントインタイムの正確性: テーブル ソースに裏付けられた集計およびColumnSelection機能の場合、将来のモデル トレーニングへのデータ漏洩を防ぐために、各行のタイムスタンプより前に利用可能なソース データのみを使用して機能がコンピュートされます。 RequestSource特徴量については、値はラベル付きDataFrameフレームの行から直接取得されます。

log_model()​

MLflowを使用して、リネージ追跡と推論中の自動特徴検索のための特徴メタデータを含むモデルをログに記録します。

Python
FeatureEngineeringClient.log_model(
model, # Trained model object
artifact_path: str, # Path to store model artifact
flavor: ModuleType, # MLflow flavor module (e.g., mlflow.sklearn)
training_set: TrainingSet, # TrainingSet used for training
registered_model_name: Optional[str], # Optional: register model in Unity Catalog
extra_pip_requirements: Optional[List[str]] = None, # Optional: Additional serving dependencies
)

flavorでは、使用するMLflowモデル フレーバーモジュール ( mlflow.sklearnやmlflow.xgboostなど) を指定します。

TrainingSetでログに記録されたモデルは、トレーニングで使用された機能へのリネージを自動的に追跡します。 トレーニングセットにRequestSourceの特徴量が含まれている場合、 RequestSource個の列が必須入力としてMLflowモデルのシグネチャに追加されます。これにより、サービス提供エンドポイントのAPIスキーマが、呼び出し元が推論時に提供する必要のあるフィールドを反映していることが保証されます。詳細については、特徴量テーブルでトレーニングするモデルをご覧ください。

FeatureViewSource について、モデルをログに記録する前に、派生特徴量とその上流特徴量を登録します。推移的な依存関係が必要とするリクエスト入力も、推論時に必要です。モデル パッケージの要件については、カスタム UDF の依存関係を参照してください。

score_batch()​

自動特徴量検索によるバッチ推論を実行する:

Python
FeatureEngineeringClient.score_batch(
model_uri: str, # URI of logged model
df: DataFrame, # DataFrame with entity keys and timestamps
) -> DataFrame

score_batch モデルとともに保存された特徴メタデータを使用して、推論用の特定時点の正しい特徴を自動的にコンピュートし、トレーニングとの一貫性を確保します。 詳細については、特徴量テーブルでトレーニングするモデルをご覧ください。

ワークフローの例​

Python
import mlflow
from databricks.feature_engineering import FeatureEngineeringClient
from sklearn.ensemble import RandomForestClassifier

fe = FeatureEngineeringClient()

# Assume features are registered in UC
# labeled_df should have columns "user_id", "transaction_time", and "is_fraud"

# 1. Create training set using Feature Views
training_set = fe.create_training_set(
df=labeled_df,
features=features,
label="is_fraud",
)

# 2. Load training data with computed features
training_df = training_set.load_df()
X = training_df.drop("is_fraud").toPandas()
y = training_df.select("is_fraud").toPandas().values.ravel()

# 3. Train model
model = RandomForestClassifier().fit(X, y)

# 4. Log model with feature metadata
with mlflow.start_run():
fe.log_model(
model=model,
artifact_path="fraud_model",
flavor=mlflow.sklearn,
training_set=training_set,
registered_model_name="main.ecommerce.fraud_model",
)

# 5. Batch scoring with automatic feature lookup
# inference_df must contain the same entity and timeseries columns
# used during training. Features are automatically computed.
predictions = fe.score_batch(
model_uri="models:/main.ecommerce.fraud_model/1",
df=inference_df,
)
predictions.display()

RequestSourceの機能を使ったトレーニング​

モデルが推論時に提供されるデータ(API呼び出しからのトランザクションの詳細など)を必要とする場合は、テーブルベースの機能と併せてRequestSource機能を使用してください。トレーニング中に、ラベル付きDataFrameからRequestSource列が抽出されます。

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
DeltaTableSource, Feature, FieldDefinition, RequestSource,
ScalarDataType, ColumnSelection,
)

fe = FeatureEngineeringClient()

# RequestSource provides transaction data at inference time
request_source = RequestSource(
schema=[
FieldDefinition(name="transaction_amount", data_type=ScalarDataType.DOUBLE),
FieldDefinition(name="vendor_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_time", data_type=ScalarDataType.DATE),
]
)

delta_source = DeltaTableSource(
catalog_name="catalog",
schema_name="schema",
table_name="vendor_data",
)

# A column selection feature from the request source (pass-through)
latest_transaction_amount = Feature(
source=request_source,
function=ColumnSelection("transaction_amount"),
name="latest_transaction_amount",
)

# A lookup feature from a delta table
vendor_category = Feature(
source=delta_source,
function=ColumnSelection("vendor_category"),
entity=["vendor_id"],
timeseries_column="transaction_time",
name="vendor_category",
)

# labels_df must contain: transaction_id, transaction_time, vendor_id,
# transaction_amount, and the label column.
ts = fe.create_training_set(
df=labels_df,
features=[latest_transaction_amount, vendor_category],
label="is_fraud",
exclude_columns=["card_id"],
)

import mlflow
from sklearn.ensemble import RandomForestClassifier

with mlflow.start_run():
training_df = ts.load_df().toPandas()
X = training_df.drop(columns=["is_fraud"])
y = training_df["is_fraud"]
model = RandomForestClassifier().fit(X, y)

# log_model() adds RequestSource columns to the MLflow model signature
fe.log_model(
model=model,
artifact_path="fraud_model",
flavor=mlflow.sklearn,
training_set=ts,
registered_model_name="catalog.schema.fraud_model",
)

CustomUDFでリクエスト値を変換する​

リクエストデータを変換するには、ColumnSelection の代わりに CustomUDF を使用します。この例では、トランザクション Logs 変換機能を使用します。

Python
log_transaction_amount = fe.get_feature(
full_name="main.ecommerce.log_transaction_amount"
)
request_df = spark.createDataFrame(
[(0.0, 0), (99.0, 1)],
"transaction_amount DOUBLE, label INT",
)

transaction_training_set = fe.create_training_set(
df=request_df,
features=[log_transaction_amount],
label="label",
)
transaction_training_set.load_df().show()

結果には、元のtransaction_amount列とlabel列に加えて、log_transaction_amountが含まれます。UDFは各行からtransaction_amountを読み取ります。計算に必要なリクエスト列は、exclude_columnsで一覧表示できません。

FeatureViewSource 特徴量を使用してトレーニングする​

FeatureViewSource では、派生特徴量を含め、CustomUDFで他の特徴量出力を消費できます。create_training_set に渡す出力を渡します。中間依存関係をリストする必要はありません。

各入力行について、Databricks は完全な依存関係グラフを解決します。

  1. エンティティキー、Timestamp、およびウィンドウ定義を使用して、テーブルバックアップされた上流特徴量をコンピュートします。利用可能な場合は、互換性のあるオフライン マテリアライゼーションを使用します。
  2. 入力 DataFrame から必要なリクエスト列を読み取り、リクエストベースの機能を評価します。
  3. 依存関係の順序で派生特徴量を評価し、各UDFが上流の結果を受け取れるようにします。

派生特徴量は、別の時間枠やポイントインタイムルックアップを追加しません。その上流の特徴量は独自の時間セマンティクスを保持します。The DataFrame must contain the entity, Timestamp, and request columns needed by those upstreams, even when only the final derived feature is requested.

たとえば、revenue_sum_7d と cost_sum_7d を組み合わせた登録済みのマージン特徴量を使用します。

Python
from databricks.feature_engineering import FeatureEngineeringClient

fe = FeatureEngineeringClient()
margin = fe.get_feature(full_name="main.ecommerce.margin")

# labeled_df contains customer_id, event_time, and label.
training_set = fe.create_training_set(
df=labeled_df,
features=[margin],
label="label",
exclude_columns=["customer_id", "event_time"],
)
training_df = training_set.load_df()

結果にはlabelとmarginが含まれます。収益およびコストの特徴量はコンピュートされますが、追加の列としては返されません。収益をトレーニングデータに含めるには、revenue = fe.get_feature(full_name="main.ecommerce.revenue_sum_7d")で取得し、features=[margin, revenue]を渡します。これはマルチレベルのチェーンにも適用されます。最終的な機能をリクエストしても、各中間出力は返されません。

リクエストに基づく特徴量、テーブルに基づく特徴量、派生特徴量を同じ features リスト内で組み合わせることができます。1つのUDF内でそれらの値を組み合わせるには、リクエスト値を特徴量として表し、FeatureViewSource内でテーブルに基づく特徴量とともに参照します。

実験を行うには、アップグラフを含め、ローカルの Feature オブジェクトを登録せずに構築します。結果を検査するには、必要に応じて label=None を指定して create_training_set を使用します。compute_features は RequestSource または FeatureViewSource をサポートしていません。

注記

クエリーあたりの Unity Catalog UDF 呼び出しの制限(limit of five Unity Catalog UDF calls per query)は、トレーニングのクエリーにも適用されます。出力として要求された特徴量だけでなく、依存関係グラフ全体で必要な UDF 呼び出しをカウントします。このクエリーの上限は、グラフの深さの上限とは別です。

カスタム UDF の依存関係​

オフライン コンピュートの場合、Unity Catalog UDF の ENVIRONMENT 句 で Python パッケージを宣言します。ノートブックにパッケージをインストールするだけでは、UDF 環境にはインストールされません。

モデルサービングの場合は、必要なパッケージを明示的にlog_modelに渡してください。UDFのENVIRONMENTも、Feature Viewを含む名前付き特徴量仕様も、これらのモデル要件を自動的には提供しません。アップストリームのUDFに必要な依存関係と、要求された特徴量出力を含めます。

transaction_training_set.load_df() で scikit-learn モデルをトレーニングした後に、同じトレーニングセットでそれをログに記録します。NumPy と互換性のあるルックアップ パッケージを含めます。

Python
import mlflow

fe.log_model(
model=model,
artifact_path="transaction_model",
flavor=mlflow.sklearn,
training_set=transaction_training_set,
registered_model_name="main.ecommerce.transaction_model",
extra_pip_requirements=[
"numpy==1.26.4",
"databricks-feature-lookup>=1.15.0",
],
)

サービングEndpointでは、RequestSource 機能のオンデマンド Unity Catalog UDF コンピュートをサポートする databricks-feature-lookup バージョン 1.15.0 以降が自動的に適用されるはずです。計算値の差異を避けるため、オフライン環境とサービング環境で UDF パッケージのバージョンを一致させてください。モデルのない Feature Serving Endpoint の場合は、代わりに create_feature_spec でパッケージを宣言してください。Python 依存関係を追加するを参照してください。

ストリーミング機能でのトレーニング​

ストリームを定義すると、DatabricksはストリームデータをDeltaテーブルに書き込む取り込みパイプラインを管理します。create_training_set はこの取り込みテーブルから読み取りを行い、DeltaTableSource からのバッチ特徴量と同様に、ラベル付けされた DataFrame に対してポイントインタイム結合を実行します。取り込み構成、バックフィル、および重複排除の詳細については、取り込みとバックフィルを参照してください。

例​

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
StreamSource,
Feature,
AggregationFunction,
Sum,
RollingWindow,
)
from datetime import timedelta

fe = FeatureEngineeringClient()

# Define a streaming feature
stream_source = StreamSource(full_name="my_catalog.my_schema.my_stream")

streaming_feature = Feature(
name="user_purchase_sum",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(
operator=Sum(input="value.amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
)

# Create training set — reads from the ingestion table
# labeled_df must contain "user_id", "event_time", and label columns.
# Entity and timeseries columns use leaf node names (not value. prefixes).
training_set = fe.create_training_set(
df=labeled_df,
features=[streaming_feature],
label="is_fraud",
)

training_df = training_set.load_df()

バッチとストリーミングの機能の混合​

バッチ特徴量とストリーミング特徴量は、同じトレーニングセットとモデルで併用できます。サービング時には、バッチ特徴量はオフラインまたはオンラインストアから参照され、ストリーミング特徴量はオンラインストアから参照されます。

Python
training_set = fe.create_training_set(
df=labeled_df,
features=[batch_feature, streaming_feature],
label="is_fraud",
)

log_model() でログに記録されたモデルは、オンラインストアから特徴量ルックアップを実行し、両方のソースタイプに対してモデルシグネチャを設定します。

サービング時に生のモデルに届くもの​

Feature Storeモデルラッパーは、列をフィルタリングしてから、生のモデルに渡します。

列のタイプ

内部モデルに到達しますか?

明示的な特徴出力( ColumnSelection 、集約)

はい

RequestSource 特徴量として宣言された列

はい

エンティティ列(ルックアップキー)

いいえ(機能として明示的に宣言されていない限り)

時系列コラム

いいえ(機能として明示的に宣言されていない限り)

列のタイプ

内部モデルに到達しますか?

明示的な特徴出力( ColumnSelection 、集約)

はい

RequestSource 特徴量として宣言された列

はい

エンティティ列(ルックアップキー)

いいえ(機能として明示的に宣言されていない限り)

時系列コラム

いいえ(機能として明示的に宣言されていない限り)