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

特徴量ビュー

備考

プレビュー

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

フィーチャービューを使用すると、データソースから特徴量を定義およびコンピュートできます。特徴量は、さまざまなソース(Deltaテーブル、Kafkaストリーム、リクエスト時データ)とコンピュテーション(時間枠集計、単純な列選択など)を使用して定義できます。このガイドでは、以下のワークフローについて説明します。

  • 特徴量開発 ワークフロー

    • create_featureを使用して、モデルのトレーニングとサービング ワークフローで使用できる Unity Catalog 特徴量オブジェクトを定義します。
    • あるいは、 Featureオブジェクトをローカルで構築し、後でregister_featureを使用してそれらをUnity Catalogに永続化することもできます。 ローカルで構築された特徴量は、登録前にcreate_training_setと組み合わせて使用できます。
  • モデルトレーニング ワークフロー

  • 特徴量のマテリアライズとサービング のワークフロー

    • create_featureで特徴量を定義するか、 get_featureを使用して特徴量を取得した後、 materialize_featuresを使用して、その特徴量または特徴量のセットをオフライン ストアにマテリアライズして効率的に再利用したり、オンライン ストアにマテリアライズしてオンラインで提供したりできます。
    • オフラインバッチトレーニングデータセットを準備するには、マテリアライズドビューとcreate_training_setを使用します。

API の詳細については、「Feature Views API リファレンス」を参照してください。

要件​

  • Serverless コンピュートまたは従来のコンピュート クラスター (Databricks Runtime 17.0 機械学習以降を実行)。

  • カスタムPythonパッケージをインストールする必要があります。ノートブックをランするたびに、次のコード行をランします。

    Python
    %pip install databricks-feature-engineering>=0.16.0
    dbutils.library.restartPython()

クイックスタートの例​

実行可能なクイックスタートノートブックについては、「ノートブックの例」を参照してください。

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
CronSchedule, DeltaTableSource, Feature, AggregationFunction,
Sum, Avg, ColumnSelection, TableTrigger,
TumblingWindow, SlidingWindow,
OfflineStoreConfig, OnlineStoreConfig,
)
from datetime import timedelta

CATALOG_NAME = "main"
SCHEMA_NAME = "feature_store"
TABLE_NAME = "transactions"

# 1. Create data source
source = DeltaTableSource(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name=TABLE_NAME,
)

# 2. Define features locally (no catalog/schema needed yet)
avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), TumblingWindow(window_duration=timedelta(days=30))),
name="avg_transaction_30d",
)

sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), SlidingWindow(window_duration=timedelta(days=7), slide_duration=timedelta(days=1))),
# name auto-generated: "amount_sum_sliding_7d_1d"
)

fe = FeatureEngineeringClient()

# 3. Explore features with compute_features
feature_df = fe.compute_features(features=[avg_feature, sum_feature])
feature_df.display()

# 4. Create training set using local features
# `labeled_df` should have columns "user_id", "transaction_time", and "target".
training_set = fe.create_training_set(
df=labeled_df,
features=[avg_feature, sum_feature],
label="target",
)
training_set.load_df().display()

# 5. Register features in Unity Catalog
avg_feature = fe.register_feature(
feature=avg_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)
sum_feature = fe.register_feature(
feature=sum_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)

# 6. Or use create_feature for a one-step define-and-register workflow
latest_amount = fe.create_feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
name="latest_amount",
)

# 7. Train model
with mlflow.start_run():
training_df = training_set.load_df()

# training code

fe.log_model(
model=model,
artifact_path="recommendation_model",
flavor=mlflow.sklearn,
training_set=training_set,
registered_model_name=f"{CATALOG_NAME}.{SCHEMA_NAME}.recommendation_model",
)

# 8. (Optional) Materialize features for serving
# Features must be registered in UC before calling materialize_features
online_config = OnlineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features_serving",
online_store_name="customer_features_store",
)

# Aggregation features support CronSchedule or TableTrigger, and support both offline and online configs
fe.materialize_features(
features=[avg_feature, sum_feature],
offline_config=OfflineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features",
),
online_config=online_config,
trigger=CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
),
)

# ColumnSelection features use TableTrigger and only support online config
fe.materialize_features(
features=[latest_amount],
online_config=online_config,
trigger=TableTrigger(),
)

「ノートブックの例」​

フィーチャービュー クイックスタート ノートブック

ストリーミング特徴量​

バッチスケジュールではなく、特徴量を継続的に更新する必要がある場合は、ストリーミング特徴量を使用します。ストリーミング特徴量とバッチ特徴量は、同じ Feature コンストラクター、集計関数、およびトレーニングとサービングのワークフローを使用します。

ストリーミング特徴量は、オフラインストアにはマテリアライズされません。トレーニングおよびバッチ推論のために、Databricks はソースから特徴量値を計算します。

ストリーミング特徴量には、次の要件があります。

  • online_configを提供する必要があります。ストリーミング特徴量はoffline_configをサポートしていません。
  • 1回のmaterialize_features呼び出しで、ストリーミング特徴量とバッチ特徴量を組み合わせることはできません。Triggerタイプごとに個別の呼び出しを行ってください。
  • transformation_sql は StreamSource でのみ使用でき、DeltaTableSource では使用できません。
  • ストリーミングマテリアライズ処理は、パイプラインの起動後に到着したレコードのみを処理し、履歴レコードのバックフィルは行いません。ローリングウィンドウ集計は、データの最初の完全なウィンドウが到着した後にのみ、完全な結果を返します。

ストリーミング特徴量ソースの選択​

鮮度要件と既存の取り込み設定に基づいてソースを選択してください:

  • 1秒未満の鮮度が優先される場合は、StreamSource を使用してください。StreamSource 特徴量は、200 ミリ秒の p99 エンドツーエンドレイテンシーを提供します。まず ストリームを設定し、次に StreamSource を使用してそれを参照します。ストリームソースは Kafka を入力としてサポートしており、トレーニング用のデータの履歴コピーとしてインジェスト Delta テーブルを自動的に保持します。
  • Delta テーブルへの低レイテンシーのインジェストパスがすでにある場合は、DeltaTableSource を使用してください。数十秒程度の鮮度が期待できます。
  • 低レイテンシーのインジェストパスがまだない場合は、Zerobus を使用して DeltaTableSource を設定します。Zerobus のインジェストには数十秒かかるため、1分未満の特徴量の鮮度が期待できます。

StreamSourceを使用してストリーミング特徴量を定義する​

StreamSourceはストリームをその3部構成の名前(catalog.schema.stream_name)で参照します。ストリームはUnity Catalogのセキュリティ保護可能なオブジェクトではありませんが、Unity Catalogスキーマにスコープされており、アクセスはストリームの取り込みテーブルによって管理されます。エンティティ、時系列、および関数定義での列の参照には、Kafkaメッセージのどの部分を読み取るかを示すために、value.またはkey.を接頭辞として付ける必要があります。ネストされたフィールドは、ドット表記(例:value.user.address.city)を使用してサポートされています。

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

client = FeatureEngineeringClient()

stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
)

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

DeltaTableSource を使用したストリーミング特徴量の定義​

DeltaTableSource 上で定義された特徴量をストリーミング特徴量としてマテリアライズするには、StreamingMode をTriggerとして materialize_features に渡します。この特徴量定義は、DeltaTableSource に裏打ちされたバッチ特徴量と同じ APIs を使用します。Delta テーブルソースは、集計および列選択機能をサポートしています。

ソース Delta テーブルでは、delta.enableChangeDataFeed=true を設定して チェンジデータフィード (CDF) を有効にする必要があります。

次の例では、Deltaテーブルソースを使用して集計特徴量を定義およびマテリアライズします。

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
DeltaTableSource,
AggregationFunction,
OnlineStoreConfig,
Sum,
RollingWindow,
StreamingMode,
)
from datetime import timedelta

client = FeatureEngineeringClient()

source = DeltaTableSource(
catalog_name="my_catalog",
schema_name="my_schema",
table_name="transactions",
)

feature = client.create_feature(
catalog_name="my_catalog",
schema_name="my_schema",
name="user_purchase_sum",
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(
operator=Sum(input="amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
)

online_config = OnlineStoreConfig(
catalog_name="my_catalog",
schema_name="my_schema",
table_name_prefix="streaming_features",
online_store_name="my_online_store",
)

client.materialize_features(
features=[feature],
online_config=online_config,
trigger=StreamingMode(),
)

Zerobus によってデータが入力される Delta テーブルを使用する​

Zerobus によってデータが投入された Delta テーブルは、ストリーミング機能ソースとして機能します。Zerobus は delta.enableChangeDataFeed=true を自動的に設定しません。ターゲット Delta テーブルをストリーミング機能ソースとして使用する前に、このプロパティを手動で設定する必要があります。

ストリーミングソースに対するフィルタ条件​

StreamSource または DeltaTableSource の集計前に、filter_condition を使用して行をフィルタリングします。

Python
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)

StreamSourceの変換​

transformation_sql key および value 構造体へのドットプレフィックス参照を使用して、集計または列選択の前に Spark SQL プロジェクションを適用します。行単位の式のみをサポートしています。サポートされている操作およびサポートされていない操作の完全なリストについては、サポートされている transformation_sql 式を参照してください。

transformation_sql を設定する場合は、dataframe_schema も指定する必要があります。これは、投影された出力の Spark StructType JSON です。このスキーマを手動で記述する代わりに、ストリームの取り込みテーブルに対して同じプロジェクションを実行して導出することもできます。これにより、key および value 構造体が公開されます。LIMIT 0 を指定することで、SELECT ステートメントはデータを読み取らずにスキーマを返します:

Python
from databricks.feature_engineering.entities import StreamSource

transformation_sql = (
"value.amount * value.conversion_rate AS converted_amount, "
"struct(value.user_id AS user_id, value.event_time AS time) AS event"
)

ingestion_table = "my_catalog.my_schema.events_ingestion"
dataframe_schema = spark.sql(
f"SELECT {transformation_sql} FROM {ingestion_table} LIMIT 0"
).schema.json()

stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
transformation_sql=transformation_sql,
dataframe_schema=dataframe_schema,
)

ストリーミングソースからの列選択​

ColumnSelection 特徴量は、集計を行わずに、各エンティティキーのソースから最新の値を渡します。トレーニング時には、特徴値はポイントインタイムの正確性を尊重します。

列選択特徴量にはTTLがありません。オンラインストアから選択した値を削除するには、ソースが選択した列に対してNull値を出力する必要があります。

Python
from databricks.feature_engineering.entities import ColumnSelection

passenger_count = Feature(
name="passenger_count",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=ColumnSelection(column="value.passenger_count"),
)

StreamSourceから入れ子フィールドにアクセスする​

StreamSource の場合、ドット表記を使用してネストされた JSON フィールドにアクセスできます(例:value.nested_field.amount)。サービング時、リクエストペイロードとレスポンスはリーフノード名を使用します(例:value.amount ではなく amount)。サービングEndpointはリーフ名を使用して値をルーティングするため、モデルまたは特徴量仕様内のすべてのエンティティ名、時系列名、および特徴量名において、リーフノード名は一意である必要があります。

ストリーミング特徴量の時間ウィンドウ​

ストリーミング特徴量は、集計に対して RollingWindow および SawtoothWindow をサポートしています。どちらも特徴量の値をリアルタイムで最新に保ちますが、効率的にコンピュートするルックバックの長さが異なります。TumblingWindow と SlidingWindow は、固定された過去の期間に対するバッチ処理用に設計されており、ストリーミングソースで使用することはできません。

  • 短いルックバックウィンドウには、RollingWindowを使用します。RollingWindow 特徴量はストリームからのデータに基づいて特徴量データをコンピュートし、データの出入りの際にサブ秒単位の精度を保証します。特徴量集計値は、ストリーム自体からリアルタイムで構築され、ローカルに保存されます。ウィンドウが完全に埋まる前に、マテリアライズを window_duration 間実行しておく必要があります。RollingWindow は、数か月または数年のウィンドウ長に対して、ソートゥース(ノコギリ波)ほど適切にスケーリングできません。
  • 長期的なルックバック期間には、SawtoothWindowを使用してください。SawtoothWindow特徴量は、ストリームから直近2日間のデータを読み込み、ストリームの取り込みテーブルからのバッチデータでこの最先端を拡張します。これにより、取り込みテーブルに少なくともフルウィンドウの期間をカバーするデータがすでに含まれているという条件で、Databricksは長期ウィンドウのストリーミング特徴量をすばやく初期化できます。SawtoothWindowは、マテリアライズの開始から約2日後にサービングの準備が整います。そのwindow_durationは2日間(強制される最小値)より長くする必要があり、Sum、Avg、Count、Min、Max、First、Last、VarPop、VarSamp、StddevPop、およびStddevSampの集計関数をサポートしています。

Databricks では、7 日を超える特徴量期間に対して SawtoothWindow を推奨しています。2日超から7日までのウィンドウでは、RollingWindow の固定長精度を重視するか、SawtoothWindow のより迅速な本番運用への移行を重視するかを決定します。

備考

ベータ版

SawtoothWindow はベータ版です。

以下の例では、StreamSourceに対してノコギリ波ウィンドウを使用した30日間のストリーミング特徴量を定義しています。

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

stream_source = StreamSource(full_name="my_catalog.my_schema.my_stream")

user_spend_30d = Feature(
name="user_spend_30d",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(
operator=Sum(input="value.amount"),
time_window=SawtoothWindow(window_duration=timedelta(days=30)),
),
)

ストリーミング特徴量例ノートブック​

ストリーミング特徴量ビュークイックスタートノートブック

モデルトレーニングと推論​

log_model()、score_batch()、およびcreate_training_set()などの特徴量ビューを使用してモデルをトレーニングし、バッチ推論を実行するには、特徴量ビューを使用したモデルのトレーニングを参照してください。

特徴量のマテリアライズ​

特徴量を定義した後、それらをオフラインまたはオンラインストアにマテリアライズすることで、トレーニングやサービングのワークフローにおいて効率的に再利用できます。Delta テーブルソースによってバックアップされる特徴量の場合、Databricks はチェンジデータフィード (CDF) を使用してソースデータをインクリメンタルに処理し、継続的なマテリアライズを効率的に維持します。特徴量をマテリアライズした後、CPU モデルサービングを使用してモデルを提供できます。詳細については、「特徴量ビューのマテリアライズ」を参照してください。

ベストプラクティス​

特徴量の命名​

  • ビジネスに不可欠な特徴量にはわかりやすい名前を使用します。
  • チーム間で一貫した命名規則に従ってください。
  • 特徴量の開発を始める際は、自動生成された名前を使用してください。

時間ウィンドウ​

  • ウィンドウ境界をビジネスサイクル(日次、週次)に合わせます。
  • 短い期間は最近の傾向を捉えますが、ノイズが発生する可能性があります。より長いウィンドウは、より安定した特徴量の分布を生成しますが、最近の行動の変化を見逃す可能性があります。ユースケースに応じて、基になるシグナルがどのくらい早く変化するかに基づいて選択します。たとえば、7日間のウィンドウは日々の変動を平滑化し、一貫したモデル入力を生成します。一方、1時間のウィンドウは行動の変化に素早く反応しますが、モデルのパフォーマンスを低下させる分散を導入する可能性があります。分布がシフトしたときにモデルの精度が低下する場合、入力の安定化のために、より長いウィンドウを使用します。
  • タンブリングウィンドウとスライディングウィンドウは、ローリング (連続) ウィンドウよりもスケーラブルです。ほとんどのユースケースでは、スライディングウィンドウから始めます。

パフォーマンス​

  • データスキャンを最小限に抑えるため、同じデータソースからの特徴量を単一のmaterialize_features呼び出しでマテリアライズします。
  • マテリアライズ時のグループ化を改善するため、同じデータソース上のフィーチャーには同じ粒度(例えば、すべての1時間またはすべての1日スライドの期間)を使用します。

エンティティ列とフィルター条件​

同じソーステーブルの機能を使用する場合は、この決定ガイドを使用します:

異なる集計レベルが必要な場合は、entity (create_feature上) を使用します:

  • 顧客レベルの特徴量 (顧客ごとに 1 行): entity=["customer_id"]
  • 顧客と販売者の特徴量 (顧客ごとに複数の行): entity=["customer_id", "merchant_id"]
  • 異なる集計レベルで同じDeltaTableSourceを共有できます :各特徴量定義で異なるentity値を指定してください

同じ集計レベルで行をフィルタリングする必要がある場合は、DeltaTableSourceでfilter_conditionを使用してください。

  • 高額トランザクションのみ : filter_condition="amount > 100" (顧客ごとに集計されます)
  • 完了した注文のみ :filter_condition="status = 'completed'" (顧客ごとに引き続き集計されます)

経験則として、 変更によってエンティティ値ごとの行数が変わる場合は、特徴量定義で異なるentity値を使用してください。同じ集計に寄与する行をフィルタリングするだけであれば、ソースでfilter_conditionを使用してください。

一般的なパターン​

顧客アナリティクス​

Python
from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow

fe = FeatureEngineeringClient()
features = [
# Recency: Number of transactions in the last day
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=1)))),

# Frequency: transaction count over the last 90 days
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=90)))),

# Monetary: total spend in the last month
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30)))),
]

トレンド分析​

Python
# Compare recent vs. historical behavior
fe = FeatureEngineeringClient()
recent_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)

historical_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7), delay=timedelta(days=7))),
)

「季節パターン」​

Python
# Same day of week, 4 weeks ago
fe = FeatureEngineeringClient()
weekly_pattern = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=1), delay=timedelta(weeks=4))),
)

制限事項​

  • create_training_set APIで使用する場合、トレーニング(ラベル付き)データセットと特徴量定義の間で、エンティティと時系列列の名前が一致している必要があります。
  • トレーニングデータセットでlabel列として使用される列名は、Featureを定義するために使用されるソーステーブルに存在してはいけません。
  • create_feature APIでは、限られた関数 ( UDAF ) がサポートされています。 サポートされている関数を参照してください。
  • エンティティ列は、DATE型またはTIMESTAMP型にすることはできません。
  • RequestSource ScalarDataTypeで定義されているスカラーデータ型(INTEGER、FLOAT、BOOLEAN、STRING、DOUBLE、LONG、TIMESTAMP、DATE、SHORT)のみをサポートします。配列、マップ、構造体などの複合型はサポートされていません。
  • RequestSource 集計関数または時間ウィンドウはサポートしていません。使用できるのはColumnSelection関数のみです。
  • エンティティ列名、時系列列名、およびリクエスト機能列名のセットは、トレーニングセットまたはサービングEndpoint内のすべてのソースでグローバルに一意である必要があります。

マテリアライズに関する制限事項については、 「制限事項」を参照してください。