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

特徴量ビューAPIリファレンス

備考

プレビュー

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

アクセス制御​

特徴量は、ガバナンス可能な Unity Catalog オブジェクトです。特徴量へのアクセスは、CREATE FEATURE、READ FEATURE、および MANAGE の Unity Catalog 権限によって制御されます。詳細については、Unity Catalog 権限のリファレンスを参照してください。

  • CREATE FEATURE : Required to create a feature in a schema.create_feature and register_feature require CREATE FEATURE on the parent schema.Following the principle of least privilege, grant CREATE FEATURE at the schema level; you can also grant it on a catalog to allow creating features in any schema in that catalog.
  • READ FEATURE : フィーチャのメタデータを読み取るために必要です。get_feature、create_training_set、および list_materialized_features では、フィーチャに対して READ FEATURE が必要です。この権限では、ソースまたはマテリアライズド出力テーブル内のフィーチャデータへのアクセスは付与されません。トレーニングやサービングのためにそのデータを読み取るには、該当するテーブルに対する SELECT も必要です。スキーマまたはカタログで付与された READ FEATURE は、そのスキーマまたはカタログに含まれる現在および将来のすべてのフィーチャに適用されます。
  • MANAGE : 特徴量のライフサイクルと権限を管理するために必要です。delete_featureによる特徴量の削除と、materialize_featuresによる特徴量のマテリアライズには、その特徴量に対するMANAGEが必要です。delete_materialized_featureを使用したマテリアライズされた特徴量の削除はMANAGEの管理対象外であり、マテリアライズされた特徴量の作成者のみが削除できます。

すべての特徴量操作には、親カタログに対するUSE CATALOGと親スキーマに対するUSE SCHEMAも必要です。MANAGEとREAD FEATUREがマテリアライズに適用される方法については、権限を参照してください。

特徴量ビューAPI​

Feature コンストラクタと register_feature()​

推奨されるアプローチは、Featureオブジェクトをローカルで構築し、register_featureを使用してUnity Catalogに永続化することです。この2段階のワークフローでは、特徴量(create_training_setを含む)を登録する前に試すことができます。

Python
Feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, RequestSource, or FeatureViewSource
function: Union[AggregationFunction, ColumnSelection, CustomUDF], # Required for all features
entity: Optional[List[str]] = None, # Required for DeltaTableSource and StreamSource
timeseries_column: Optional[str] = None, # Required for DeltaTableSource and StreamSource
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
)

FeatureEngineeringClient.register_feature() ローカルで構築された Feature を Unity Catalog に登録する。

Python
FeatureEngineeringClient.register_feature(
feature: Feature, # Required: A Feature instance (not already registered)
catalog_name: str, # Required: UC catalog name
schema_name: str, # Required: UC schema name
) -> Feature
Python
from databricks.feature_engineering.entities import Feature, DeltaTableSource, AggregationFunction, Sum, RollingWindow
from datetime import timedelta

# Step 1: Construct the feature locally
feature = Feature(
source=DeltaTableSource(catalog_name="main", schema_name="store", table_name="transactions"),
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)

# Step 2: Register in Unity Catalog
fe = FeatureEngineeringClient()
registered_feature = fe.register_feature(
feature=feature,
catalog_name="main",
schema_name="store",
)

create_feature()​

FeatureEngineeringClient.create_feature() 単一のステップで、Unity Catalog内の機能を検証、構築し、すぐに登録します。最初にローカルで機能のエクスペリメントをする必要がない場合、これを使用します。

Python
FeatureEngineeringClient.create_feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, RequestSource, or FeatureViewSource
function: Union[AggregationFunction, ColumnSelection, CustomUDF], # Required for all features
catalog_name: str, # Required: The catalog name for the feature
schema_name: str, # Required: The schema name for the feature
entity: Optional[List[str]] = None, # Required for DeltaTableSource and StreamSource
timeseries_column: Optional[str] = None, # Required for DeltaTableSource and StreamSource
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
) -> Feature

パラメーター:

  • source:特徴量計算で使用されるデータソース(DeltaTableSource、StreamSource、RequestSource、または FeatureViewSource)。
  • function: オペレーターと時間枠をバンドルする AggregationFunction、パススルー特徴量の ColumnSelection("column_name")、または行ごとの変換の CustomUDF。互換性のあるソースタイプについては、サポートされている関数を参照してください。
  • catalog_name特徴量に対するUnity Catalogのカタログ名。
  • schema_name:特徴量のUnity Catalogスキーマ名です。
  • entity集計キーまたは検索キー(プライマリキー)を定義する列名のリスト。DeltaTableSource および StreamSource に必要です。たとえば、["user_id"] はユーザーごとに集計または検索を行います。RequestSource および FeatureViewSource では省略します。
  • timeseries_column:時間ウィンドウ集計または最新値の選択に使用される Timestamp 列。DeltaTableSource および StreamSource に必要です。RequestSource および FeatureViewSource では省略します。
  • name:任意の機能名。省略した場合、入力列、関数、およびウィンドウから自動生成されます(例:amount_avg_rolling_7d)。
  • description:特徴量のオプションの説明。

戻り値: 検証済みの特徴量インスタンス

発生: いずれかの検証が失敗した場合、ValueErrorが発生します。

delete_feature()​

Unity Catalogから、その完全修飾名で特徴量を削除します。

Python
FeatureEngineeringClient.delete_feature(
full_name: str, # Required: '<catalog>.<schema>.<feature_name>'
) -> None
Python
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")

機能を削除する前に、それを参照しているモデルやfeature specを削除または更新してください。マテリアライズされた特徴量が残っている間は、機能を削除できません。最初にマテリアライズされた特徴量を削除してから、機能を削除してください。How to delete a materialized featureを参照してください。

自動生成された名前​

nameが省略されている場合、名前は自動的に生成されます。生成される名前はパターン:{column}_{function}_{window}に従います。例えば:

  • price_avg_rolling_1h (1時間あたりの平均価格)
  • transaction_count_rolling_30d_1d (イベントの Timestamp からの 1 日遅延を伴う 30 日間のトランザクション数)

サポートされている関数​

集計関数​

注記

集計関数は、時間ウィンドウに記載されているように、時間ウィンドウと共に AggregationFunction でラップされています。各関数は、集計するソース列を指定する input パラメーターを受け取ります。

関数

説明

使用例

Sum(input="column")

値の合計

ユーザーごとの1日のアプリ使用量 (単位: 分)

Avg(input="column")

平均値

平均トランザクション額

Count(input="column")

レコード数

ユーザーあたりのログイン数

Min(input="column")

最小値

ウェアラブルデバイスによって記録された最低心拍数

Max(input="column")

最大値

セッションあたりの最大トランザクション量

StddevPop(input="column")

母集団標準偏差

すべての顧客における日次取引額の変動

StddevSamp(input="column")

サンプル標準偏差

広告キャンペーンのクリック率の変動性

VarPop(input="column")

母集団分散

工場におけるIoTデバイスのセンサー読み取り値の分散

VarSamp(input="column")

サンプルバリアンス

サンプリングされたグループ全体での映画の評価の広がり

ApproxCountDistinct(input="column", relativeSD=0.05)

おおよそのユニークカウント

購入されたアイテムの個別カウント

ApproxPercentile(input="column", percentile=0.95, accuracy=100)

近似パーセンタイル

p95 応答レイテンシ

First(input="column")

最初の値

最初のログインTimestamp

Last(input="column")

最後の値

最新の購入額

FirstN(input="column", n=3)

配列としての最初の n 件の値

セッションで閲覧された最初の3つの製品

LastN(input="column", n=3)

配列としての最後の n 件の値

最新のサポートケースステータス3件

FirstDistinct(input="column", n=3)

配列としての最初の n 件の個別の値

閲覧された最初の3つの個別の製品カテゴリ

LastDistinct(input="column", n=3)

配列としての最後の n 件の個別の値

最新の3件の個別の加盟店カテゴリ

関数

説明

使用例

Sum(input="column")

値の合計

ユーザーごとの1日のアプリ使用量 (単位: 分)

Avg(input="column")

平均値

平均トランザクション額

Count(input="column")

レコード数

ユーザーあたりのログイン数

Min(input="column")

最小値

ウェアラブルデバイスによって記録された最低心拍数

Max(input="column")

最大値

セッションあたりの最大トランザクション量

StddevPop(input="column")

母集団標準偏差

すべての顧客における日次取引額の変動

StddevSamp(input="column")

サンプル標準偏差

広告キャンペーンのクリック率の変動性

VarPop(input="column")

母集団分散

工場におけるIoTデバイスのセンサー読み取り値の分散

VarSamp(input="column")

サンプルバリアンス

サンプリングされたグループ全体での映画の評価の広がり

ApproxCountDistinct(input="column", relativeSD=0.05)

おおよそのユニークカウント

購入されたアイテムの個別カウント

ApproxPercentile(input="column", percentile=0.95, accuracy=100)

近似パーセンタイル

p95 応答レイテンシ

First(input="column")

最初の値

最初のログインTimestamp

Last(input="column")

最後の値

最新の購入額

FirstN(input="column", n=3)

配列としての最初の n 件の値

セッションで閲覧された最初の3つの製品

LastN(input="column", n=3)

配列としての最後の n 件の値

最新のサポートケースステータス3件

FirstDistinct(input="column", n=3)

配列としての最初の n 件の個別の値

閲覧された最初の3つの個別の製品カテゴリ

LastDistinct(input="column", n=3)

配列としての最後の n 件の個別の値

最新の3件の個別の加盟店カテゴリ

注記

First、Last、FirstN、LastN、FirstDistinct、および LastDistinct は、defaultでNull値を含みます。Null値をスキップするには、Nullである入力列を明示的に除外する filter_condition を追加します。

FirstN、LastN、FirstDistinct、および LastDistinct は、特徴量の timeseries_column を使用して入力行を順序付けし、最大 n 個の値を含む配列を返します。n パラメーターには正の整数を指定する必要があります。FirstN および FirstDistinct は、最も古いものから最新のものへと値を選択します。LastN および LastDistinct は、最新から最も古い順に値を選択し、選択した値をTimestamp順に返します。FirstDistinct および LastDistinct は、その方向に値を選択しながら重複する値を削除します。

例えば、エンティティのソース行が event_time によって ["A", "A", "B", "C", "B", "B"] として順序付けられている場合、以下の関数は次を返します。

関数

結果

FirstN(input="event_type", n=3)

["A", "A", "B"]

LastN(input="event_type", n=3)

["C", "B", "B"]

FirstDistinct(input="event_type", n=3)

["A", "B", "C"]

LastDistinct(input="event_type", n=3)

["A", "C", "B"]

関数

結果

FirstN(input="event_type", n=3)

["A", "A", "B"]

LastN(input="event_type", n=3)

["C", "B", "B"]

FirstDistinct(input="event_type", n=3)

["A", "B", "C"]

LastDistinct(input="event_type", n=3)

["A", "C", "B"]

FirstN、LastN、FirstDistinct、および LastDistinct には、databricks-feature-engineering バージョン 0.17.0 以降が必要です。

カスタムUDF​

CustomUDF 登録済みの Unity Catalog Python ユーザー定義関数 (UDF) を各行に適用します。これを使用して、リクエスト入力を変換するか、特徴量を組み合わせます。行を集計したり、時間枠を定義したりしません。

Python
CustomUDF(
function_name="main.ecommerce.log_amount_udf",
input_bindings={"amount": "transaction_amount"},
)

input_bindings 各 UDF パラメーター名をインプットにマップします。RequestSourceの場合、入力はソース列名です。FeatureViewSourceについては、上流の機能参照です。default値を持つパラメーターを含め、すべての UDF パラメーターをバインドします。暗黙的な数値キャストを行わずに、入力タイプが UDF パラメーターのタイプと完全に一致している必要があります。スカラー入力と戻り値のタイプを使用します。

ソース

挙動

RequestSource

トレーニング DataFrame または推論リクエストの列を変換します。

FeatureViewSource

上流の特徴量値を結合します。「FeatureViewSource」を参照してください。

ソース

挙動

RequestSource

トレーニング DataFrame または推論リクエストの列を変換します。

FeatureViewSource

上流の特徴量値を結合します。「FeatureViewSource」を参照してください。

Delta をバックアップとした CustomUDF の特徴量は、マテリアライズまたはオンラインサービングを行うことができません。トレーニングとサービングのためにテーブルをバックアップとした特徴量値を変換するには、Delta をバックアップとした集計または列選択の特徴量を定義し、FeatureViewSource を介してそれを参照します。

CustomUDF StreamSourceではサポートされていません。ストリーミング特徴量の出力を変換するには、FeatureViewSource を介してその特徴量を参照します。

CustomUDF RequestSource を使用するには、databricks-feature-engineering バージョン 0.17.0 以降が必要です。

CustomUDFを使用するには、UDFに対するEXECUTE権限、その親カタログに対するUSE CATALOG権限、およびその親スキーマに対するUSE SCHEMA権限が必要です。

次の例では、NumPy を使用して log(1 + amount) をコンピュートし、大規模なトランザクション金額のスケールを削減します。カスタム UDF 依存関係を有効にして、Serverless コンピュート上で実行します。main.ecommerce スキーマが存在している必要があります。

Python
spark.sql("""
CREATE OR REPLACE FUNCTION main.ecommerce.log_amount_udf(amount DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
ENVIRONMENT (
dependencies = '["numpy==1.26.4"]',
environment_version = '5'
)
AS $$
import numpy as np

if amount is None or not np.isfinite(amount) or amount < 0:
return None
return float(np.log1p(amount))
$$
""")

リクエスト列 transaction_amount を UDF パラメーター amount にバインドする特徴量を登録する:

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
CustomUDF, FieldDefinition, RequestSource, ScalarDataType,
)

fe = FeatureEngineeringClient()

log_transaction_amount = fe.create_feature(
source=RequestSource(
schema=[
FieldDefinition(name="transaction_amount", data_type=ScalarDataType.DOUBLE),
]
),
function=CustomUDF(
function_name="main.ecommerce.log_amount_udf",
input_bindings={"amount": "transaction_amount"},
),
catalog_name="main",
schema_name="ecommerce",
name="log_transaction_amount",
)

UDF の ENVIRONMENT は、オフライン計算の依存関係を構成します。オンラインサービングの場合は、create_feature_spec(extra_pip_requirements=...) または log_model(extra_pip_requirements=...) でもパッケージを宣言してください。これらは UDF から自動的にはコピーされません。Feature Serving の依存関係およびモデルの依存関係をご覧ください。

CustomUDF 特徴量をマテリアライズすることはできません。リクエストバックおよび特徴量バックの UDF は、トレーニングおよびサービング中にオンデマンドで実行されます。依存関係チェーン内の各 UDF はコンピュートを追加するため、関数とチェーンは小さく保ってください。UDF は、オフラインで None またはオンラインで NaN になる可能性がある、欠落した入力を処理する必要があります。

欠損値の処理に関するガイダンスについては、欠損している特徴量の値を処理する方法を参照してください。

ColumnSelection(パススルー)​

ColumnSelection 集計を適用せずにソースから単一の列を選択します。それは直接 function パラメーターにラップされています(AggregationFunction の中ではありません)。戻り値の型はソーススキーマから推論されます。

関数

説明

使用例

ColumnSelection("col")

列の最新値(集計なし)

最新のベンダーカテゴリ、リクエストフィールドのパススルー

関数

説明

使用例

ColumnSelection("col")

列の最新値(集計なし)

最新のベンダーカテゴリ、リクエストフィールドのパススルー

ColumnSelection は次のデータソースをサポートしています。

  • DeltaTableSource :時点結合によりエンティティキーごとに最新の値を返します(ルックバックウィンドウ集計なし)。
  • StreamSource : ストリームからエンティティキーごとの最新の値を返します(ルックバック期間の集計は行われません)。
  • RequestSource ** **:推論時に提供された値(またはトレーニング時にラベル付けされたDataFrameから抽出された値)を渡します。

DeltaTableSourceの場合、ColumnSelectionの特徴量は、最新値の選択の前に適用されるfilter_conditionおよびtransformation_sqlをサポートします(集計特徴量と同じ)。

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

delta_source = DeltaTableSource(
catalog_name="main", schema_name="feature_store", table_name="transactions",
)

request_source = RequestSource(
schema=[
FieldDefinition(name="session_duration", data_type=ScalarDataType.DOUBLE),
]
)

# ColumnSelection from a Delta table
latest_amount = Feature(
source=delta_source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
name="latest_transaction_amount",
)

# ColumnSelection from a RequestSource
session_feature = Feature(
source=request_source,
function=ColumnSelection("session_duration"),
name="session_duration",
)

例:集計と列選択機能​

次の例では、同じデータソースに対して定義された機能を示します。

Python
from databricks.feature_engineering.entities import (
AggregationFunction, Feature, Sum, Avg, ApproxCountDistinct,
ColumnSelection, RollingWindow,
)
from datetime import timedelta

window = RollingWindow(window_duration=timedelta(days=7))

sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Sum(input="amount"), window),
)

avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Avg(input="amount"), window),
)

distinct_count = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(ApproxCountDistinct(input="product_id", relativeSD=0.01), window),
)

# Column selection (no aggregation, no time window)
latest_amount = Feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="event_time",
name="latest_amount",
)

フィルター条件を持つ機能​

filter_condition パラメーターを使用すると、集計の計算や最新の列値の選択を行う 前に 、ソーステーブルから行をフィルター処理できます。これは、データのグループ化および集計の前に適用される SQL の WHERE 句として機能します。

注記

集計機能の場合、filter_conditionは、GROUP BYの前に適用されるSQL WHERE句のように、集計前に行をフィルター処理します。これにより、粒度が変更されることはありません。粒度は常に、特徴量定義のentityによって定義されます。

フィルターは、特徴量コンピュテーションに必要なデータのスーパーセットを含む大規模なソーステーブルを扱う際に役立ち、これらのテーブル上に個別のビューを作成する必要性を最小限に抑えます。

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

# Source with filter applied at the source level
high_value_transactions = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="transactions",
filter_condition="amount > 100", # Only transactions over $100
)

high_value_sales = Feature(
source=high_value_transactions,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30))),
)

# Multiple conditions
completed_orders_source = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="orders",
filter_condition="status = 'completed' AND payment_method = 'credit_card'",
)

completed_orders = Feature(
source=completed_orders_source,
entity=["user_id"],
timeseries_column="order_time",
function=AggregationFunction(Count(input="order_id"), RollingWindow(window_duration=timedelta(days=7))),
)

# Filter on a StreamSource
from databricks.feature_engineering.entities import StreamSource

purchase_stream = StreamSource(
full_name="main.ecommerce.transactions_stream",
filter_condition="value.event_type = 'purchase'",
)

purchase_total = Feature(
source=purchase_stream,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(Sum(input="value.amount"), RollingWindow(window_duration=timedelta(hours=1))),
)

データソース​

DeltaTableSource​

DeltaTableSource ソーステーブルから特徴量がどのように計算されるかを定義するために使用される、一時的なPythonオブジェクトです。新しいテーブルは作成されません。データの読み取りと特徴量の集計のための設定を指定します。

Python
DeltaTableSource(
catalog_name: str, # Required: Catalog name
schema_name: str, # Required: Schema name
table_name: str, # Required: Table name
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause to filter source data
transformation_sql: Optional[str] = None, # Optional: SQL SELECT expression for column transformations
dataframe_schema: Optional[str] = None, # Required if transformation_sql is set: schema of the resulting DataFrame
lateness: Optional[SourceLateness] = None, # Optional: Expected source settling time
)

パラメーター:

  • catalog_name、schema_name、table_name:Unity CatalogのソースDeltaテーブルを識別します。
  • filter_condition:集計または列の選択の前に適用されるSQL WHERE句。例:"status = 'completed'"。
  • transformation_sqlソーステーブルに適用される SQL SELECT 式。集計や列の選択の前に、これを使用して列の名前の変更、型のキャスト、または派生列のコンピュートを行います。省略した場合は、すべての列が選択されます(*)。例: "user_id, CAST(amount AS DOUBLE) AS amount, event_time"。
  • dataframe_schema:変換後の結果のDataFrameのスキーマ。Spark StructType JSON形式(df.schema.json()から)。 transformation_sqlが提供されている場合、必須です。 これは、変換によって生成される列名と型をシステムに伝えます。
  • lateness:イベント時間において、ソースが完全に揃うまでに通常かかる時間を記述する SourceLateness オブジェクト。省略した場合、ソースは即座に完了したと見なされます。

filter_condition と transformation_sql の両方が設定されている場合、結果のクエリーは次のようになります:SELECT {transformation_sql} FROM {table} WHERE {filter_condition}。

SourceLateness.settling_delay は、オンラインのマテリアライズベーションに影響を与える一貫したETLの遅延をトレーニング中にシミュレートするための推奨される方法です。Databricks は、トレーニングデータがオンラインでまだ転送中であったデータを使用しないように、該当するトレーニング評価時間をこの期間だけ後方にシフトします。マテリアライズベーション中、Databricksは完了したウィンドウをパブリッシュする前に同じ期間待機し、その間の期間中は最後に完了したウィンドウを提供します。

たとえば、深夜0時が協定世界時(UTC)の07:00に対応するローカルタイムゾーンにおいて、日次ETLジョブが深夜0時の8時間後に完了すると仮定します。セトリング遅延を 8 時間、ウィンドウオフセットを 7 時間に設定します。

Python
from datetime import timedelta
from databricks.feature_engineering.entities import (
DeltaTableSource,
SourceLateness,
TumblingWindow,
)

source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
lateness=SourceLateness(settling_delay=timedelta(hours=8)),
)

window = TumblingWindow(
window_duration=timedelta(days=1),
offset=timedelta(hours=7),
)
注記

timeseries_column の型は TimestampType または TimestampNTZType である必要があります。DateType は時系列ではサポートされていません。最初にカラムを TimestampType にキャストしてください(例:transformation_sql を使用)。

例:列の変換に transformation_sql を使用する

Python
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="raw_events",
transformation_sql="user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time",
filter_condition="event_type = 'purchase'",
dataframe_schema=spark.sql(
"SELECT user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time FROM main.analytics.raw_events LIMIT 0"
).schema.json(),
)

例:PySpark DataFrameからtransformation_sqlとdataframe_schemaを導出する

変換をPySparkクエリーとして記述し、その後、結果のDataFrameからスキーマを抽出できます。

Python
df = spark.sql(f"""
SELECT user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time
FROM main.analytics.events
WHERE event_date >= date_sub(current_date(), 7)
LIMIT 0
""")

# Use df.schema.json() as the dataframe_schema
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
transformation_sql="user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time",
filter_condition="event_date >= date_sub(current_date(), 7)",
dataframe_schema=df.schema.json(),
)

サポートされているtransformation_sql式​

DeltaTableSource および StreamSource 上の transformation_sql にも同じルールが適用されます。

transformation_sql 任意の行単位の式をサポートします。操作は各行に対して独立して評価されます。これらは行数やソースとの1対1の対応関係を変更しません。行単位の式には、列の名前変更、キャスト、算術演算などが含まれます。

SUM()やCOUNT()のような集計など、形状や行数を変更する操作はサポートされていません。代わりに特徴量定義でAggregationFunctionを使用してください。

DeltaTableSource.from_sql()​

便宜上、SQLクエリーからDeltaTableSourceを作成できます。このメソッドはクエリーを解析し、テーブル名、transformation_sql、およびfilter_conditionを自動的に抽出します。

Python
DeltaTableSource.from_sql(
sql: str, # Required: SQL SELECT query
spark: SparkSession, # Required: active SparkSession (for schema inference)
) -> DeltaTableSource

シンプルなSELECT ... FROM ... [WHERE ...]クエリーのみがサポートされています。複雑なSQL(JOIN、サブクエリー、CTE、UNION)は拒否されます。複雑なクエリーの場合は、transformation_sqlとfilter_conditionでDeltaTableSourceを直接構築します。

Python
from databricks.feature_engineering.entities import (
AggregationFunction,
DeltaTableSource,
Feature,
Sum,
TumblingWindow,
)

source = DeltaTableSource.from_sql(
spark=spark,
sql=f"SELECT customer_id, event_ts, amount * 2 AS doubled_amount, amount FROM {CATALOG}.{SCHEMA}.{TABLE}",
)

feature = Feature(
source=source,
function=AggregationFunction(Sum(input="doubled_amount"), time_window=TumblingWindow(window_duration=timedelta(days=7))),
entity=["customer_id"], timeseries_column="event_ts",
)

で反復 to_dataframe()​

特徴量コンピュテーションに使用されるデータをプレビューするには、source.to_dataframe()を使用します。これは、期待される結果が得られるまでfilter_conditionとtransformation_sqlを反復するのに役立ちます。

Python
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
filter_condition="event_type = 'purchase'",
)

# Preview the filtered source data
source.to_dataframe().display()

エンティティの理解​

エンティティ列は、特徴量の集計レベルを定義します。これらはFeatureの定義で指定され、DeltaTableSourceでは指定されません。エンティティが決定するもの:

  • データのグループ化方法 :特徴量は、エンティティ値の一意の組み合わせごとに集計されます(SQLのGROUP BYに類似)。
  • 主キー構造 :各一意のエンティティの組み合わせは、コンピュートされた特徴の1行になります。

例:顧客レベルの機能

次のコードは、顧客レベルで特徴量を集計します(顧客ごとに 1 行)。

Python
from databricks.feature_engineering.entities import DeltaTableSource

source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="user_events",
)

Feature(
source=source,
entity=["user_id"], # Features aggregated per user
timeseries_column="event_time", # Timestamp for time windows
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)

例:顧客・ストア レベルの特徴量

より詳細なレベルで特徴を集約するには(顧客と店舗の組み合わせごとに1行)、複数のエンティティ列を使用します:

Python
source = DeltaTableSource(
catalog_name="main",
schema_name="retail",
table_name="transactions",
)

Feature(
source=source,
entity=["user_id", "store_id"], # Features aggregated per user-store pair
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)

集計の異なるレベル(例えば、顧客レベルや顧客ストアレベル)で特徴量が必要な場合、特徴量定義で異なるentity値を使用します。異なるエンティティ構成を持つ特徴量間で同じDeltaTableSourceを共有できます。

StreamSource​

StreamSource ストリームを参照します。ストリームには、ストリーミングソースの接続、認証、スキーマ、取り込みの設定が含まれています。Kafkaの場合、特徴量定義内のカラム参照は、メッセージのどの部分を読み取るかを示すため、value.またはkey.をプレフィックスとして付ける必要があります。

Python
StreamSource(
full_name: str, # Required: Three-part Stream name (catalog.schema.stream)
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause applied before aggregation
transformation_sql: Optional[str] = None, # Optional: SQL SELECT expression for column transformations
dataframe_schema: Optional[str] = None, # Required if transformation_sql is set: schema of the resulting DataFrame
lateness: Optional[SourceLateness] = None, # Optional: Expected source settling time
)

パラメーター:

  • full_name:ストリームの完全な3部構成の名前(例:"my_catalog.my_schema.my_stream")。
  • filter_condition (オプション):集計前にストリームデータに適用されるSQL WHERE句で、ドットプレフィックス付きの列参照(例えば、"value.event_type = 'purchase'")を使用します。
  • transformation_sql (オプション): 集計または列選択の前に適用されるSQL SELECT式。keyおよびvalue構造体へのドットプレフィックス参照を使用します。DeltaTableSourceと同じ行単位の式をサポートしています。省略した場合、ソースはすべての列(*)を使用します。
  • dataframe_schema:投影された出力のSpark StructType JSONスキーマ。transformation_sqlを設定する場合に必要です。
  • lateness:ストリームがイベント時間内で完了するまでに通常かかる時間を記述する SourceLateness オブジェクト。SourceLateness.settling_delayを参照してください。
Python
from databricks.feature_engineering.entities import StreamSource

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

keyおよびvalue構造体を公開するストリームの取り込みテーブルに対してプロジェクションを実行し、dataframe_schemaを導出します。

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

RequestSource​

RequestSource 事前に実体化されたテーブルから検索されるのではなく、リクエストペイロード内の推論時に提供されるデータのスキーマを定義します。トレーニング中に、これらの列はcreate_training_setに渡されたラベル付きDataFrameから抽出されます。モデルサービング中に、呼び出し元はそれらをHTTPリクエストペイロードに含める必要があります。

RequestSource CustomUDF または ColumnSelection フィーチャー ビュー関数で使用できます。集計関数や時間枠はサポートされていません。

スキーマの定義​

スキーマをFieldDefinitionオブジェクトのリストとして定義します。各オブジェクトは、列名とScalarDataTypeを指定します。

Python
from databricks.feature_engineering.entities import (
FieldDefinition, RequestSource, ScalarDataType,
)

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

サポートされているデータ型​

RequestSource ScalarDataTypeで定義されているスカラー型をサポートします:INTEGER、FLOAT、BOOLEAN、STRING、DOUBLE、LONG、TIMESTAMP、DATE、SHORT。配列、マップ、構造体などの複合型はサポートされていません。

リクエストデータがどのようにハイドレートされるか​

コンテキスト

挙動

**トレーニング**(create_training_set )

列はラベル付きDataFrameから抽出されます。型は宣言されたスキーマに対して検証されます。不一致があるとエラーが発生します(暗黙的なキャストなし)。

サービング (モデルEndpoint)

HTTPリクエストでは、dataframe_recordsまたはdataframe_splitから列が取得されます。JSON値は、宣言された型にキャストされます(例:JSON数値 → DOUBLE)。

コンテキスト

挙動

**トレーニング**(create_training_set )

列はラベル付きDataFrameから抽出されます。型は宣言されたスキーマに対して検証されます。不一致があるとエラーが発生します(暗黙的なキャストなし)。

サービング (モデルEndpoint)

HTTPリクエストでは、dataframe_recordsまたはdataframe_splitから列が取得されます。JSON値は、宣言された型にキャストされます(例:JSON数値 → DOUBLE)。

モデルシグネチャ​

RequestSource特徴を含むトレーニングセットでlog_modelを使用してモデルがログに記録されると、必要な入力としてRequestSource列がMLflowモデル署名に追加されます。これは、サービングEndpointのAPIスキーマが、呼び出し元が推論時に提供する必要があるフィールドを反映していることを意味します。

FeatureViewSource​

FeatureViewSource 他の特徴量ビューの出力を CustomUDF の入力として使用します。特徴量を連鎖させることで、有向非巡回グラフ(DAG)が作成されます。たとえば、マージン特徴量は収益とコストの集計を組み合わせることができ、別の特徴量でそのマージンを変換できます。

FeatureViewSourceにはdatabricks-feature-engineeringのバージョン0.18.0以降を使用してください。

features には、機能名の文字列ではなく、Feature オブジェクトのリストを渡します。get_feature を使用して登録済みの特徴量を取得します。input_bindings で、登録された各特徴量の full_name を使用します。ローカルの未登録の特徴量の場合は、代わりにその name を使用します。

次の例では、2つの登録済み機能(revenue_sum_7dおよびcost_sum_7d)を想定しています。これらはDOUBLEによってcustomer_idの値目を返し、ポイントインタイム・コンピュートにevent_timeを使用します:

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import CustomUDF, FeatureViewSource

fe = FeatureEngineeringClient()
revenue = fe.get_feature(full_name="main.ecommerce.revenue_sum_7d")
cost = fe.get_feature(full_name="main.ecommerce.cost_sum_7d")

spark.sql("""
CREATE OR REPLACE FUNCTION main.ecommerce.margin_udf(revenue DOUBLE, cost DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
import math

if revenue is None or cost is None:
return None
if not math.isfinite(revenue) or not math.isfinite(cost) or revenue <= 0:
return None
return (revenue - cost) / revenue
$$
""")

margin = fe.create_feature(
source=FeatureViewSource(features=[revenue, cost]),
function=CustomUDF(
function_name="main.ecommerce.margin_udf",
input_bindings={"revenue": revenue.full_name, "cost": cost.full_name},
),
catalog_name="main",
schema_name="ecommerce",
name="margin",
)

次の制限が適用されます:

  • 関数としてサポートされているのは CustomUDF のみです。アップストリームの特徴量は、集計、列の選択、またはその他の CustomUDF の特徴量にすることができます。
  • 派生特徴量の entity と timeseries_column を省略します。各上流の特徴量は、それぞれ独自のエンティティ、Timestamp、ウィンドウ定義を保持します。
  • 特徴量には 1 つのソースがあります。リクエスト値をテーブルでバックアップされた特徴量と組み合わせるには、RequestSource 特徴量を定義し、FeatureViewSource を介して両方を参照します。
  • 宣言されたすべての上流特徴量を input_bindings で使用する必要があります。循環は許可されません。
  • 派生特徴量を登録する前に、上流の特徴量を登録します。実験のために、ローカルの未登録グラフを create_training_set と共に使用できます。
  • トレーニングまたはサービングを行うには、派生した特徴量とその推移的アップストリーム特徴量に対する READ FEATURE または MANAGE 権限が必要です。カタログやスキーマの間でも、登録とサービングのためにグラフ全体で固有の特徴量名を使用してください。
  • 1つの特徴量は、最大 20個の直接の上流特徴量を参照できます。登録済みのグラフでは、ベース特徴量を含めて、依存関係パスに沿って最大 5個の特徴量の深さがサポートされます。
  • FeatureViewSource compute_features を使用して特徴量をマテリアライズまたは評価することはできません。create_training_set を使用してそれらをオフラインで評価します。オンラインサービングを行うには、代わりにサポートされているテーブルバックアップ形式のアップストリーム特徴量をマテリアライズしてください。

依存関係の評価と出力の選択については、Train with FeatureViewSource features を参照してください。デプロイについては、Serve derived features を参照してください。

トレーニングと推論API​

create_training_set およびscore_batchが、ソースデータから特定の時点に正しい特徴量値をオンデマンドでコンピュートします。オフラインマテリアライゼーションをサポートする機能(Deltaテーブルソースでのスライディングウィンドウ集計など)の場合、まず機能をオフラインストアにマテリアライズすると、両方の操作のパフォーマンスが向上します。マテリアライズされたオフラインの特徴量が利用可能な場合、操作はソースから特徴量値を再計算する代わりに、事前に計算されたオフラインデータを読み取ります。オフラインストアに特徴量をマテリアライズするには、特徴量ビューをマテリアライズするを参照してください。

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

log_model()​

Logs a model with feature metadata for リネージ tracking and automatic feature lookup during inference.詳細については、フィーチャービューでモデルをトレーニングするを参照してください。

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
)

score_batch()​

自動の特徴量ルックアップを使用して、オフラインバッチ推論を実行します。モデルと一緒に保存されている特徴量メタデータを使用して、時点に合った特徴量をコンピュートし、トレーニングとの一貫性を確保します。

Python
FeatureEngineeringClient.score_batch(
model_uri: str, # URI of logged model (e.g., "models:/catalog.schema.model/1")
df: DataFrame, # DataFrame with entity keys and timestamps
) -> DataFrame

入力 DataFrame には、トレーニング中に使用されたエンティティ列と時系列列が含まれている必要があります。特徴はソースデータから自動的にコンピュートされます。

Python
fe = FeatureEngineeringClient()

# Batch scoring with automatic feature lookup
predictions = fe.score_batch(
model_uri="models:/main.ecommerce.fraud_model/1",
df=inference_df,
)
predictions.display()

時間ウィンドウ​

フィーチャービューは、時間ウィンドウベースの集計におけるルックバック動作を制御するための4つのウィンドウタイプをサポートしています。利用可能なウィンドウタイプは、特徴量のソースによって異なります。ストリーミングソースの特徴量ではローリングウィンドウとソートゥースウィンドウを使用でき、バッチソースの特徴量ではローリングウィンドウ、タンブリングウィンドウ、スライディングウィンドウを使用できます。

  • ローリングウィンドウはイベント時間から遡って参照します。期間と遅延は明示的に定義されています。
  • タンブリングウィンドウは、固定された重複しない時間ウィンドウです。各データポイントは1つのウィンドウにのみ属します。
  • スライディングウィンドウは、設定可能なスライド間隔を持つ、重複するローリング時間ウィンドウです。
  • Sawtooth ウィンドウは、バッチとストリーミングのハイブリッドパスを使用して、ストリーミングソース全体で長いルックバックウィンドウを最新の状態に保ちます。Sawtooth ウィンドウを参照してください。

次の図は、タンブリング、スライディング、ローリング、および Sawtooth ウィンドウのタイプを示しています。

タンブリング、スライディング、ローリング、および Sawtooth ルックバックウィンドウ。

時間ウィンドウのタイミング​

delay を使用して、過去の分析時点のウィンドウを評価します。たとえば、7日間の遅延がある30日間のウィンドウでは、評価時点の1週間前の時点として30日間の値がコンピュートされます。delay は、ソースの到着時間とは独立しています。ソースデータの到着にかかる時間をモデル化するには、代わりに SourceLateness.settling_delay を構成します。

両方の設定が存在する場合、それらは合成されます。Databricksは、ソースのセトリング遅延の経過後にウィンドウが完了したとみなし、分析遅延を使用してそれを評価します。

offsetを使用して、固定ウィンドウ境界の配置を変更します。default, タンブリングウィンドウとスライディングウィンドウは UTC の深夜 0 時に揃えられます。たとえば、22 時間のオフセットにより、1 日の境界が 22:00 UTC に合わせられます。ローカルタイムゾーンで境界を近似するには、UTC に対する静的オフセットを構成します。オフセットは、夏時間の調整、評価時間のシフト、または到着遅延データのモデル化を行いません。

次の表は、これらのフィールドのサポート状況をまとめたものです。

フィールド

サポートされているウィンドウ

制約

delay

ローリング、タンブリング、スライディング

負の値ではない必要があります。 datetime.timedelta

offset

タンブリングおよびスライディング

0以上であり、期間*よりも短い必要があります。

SourceLateness.settling_delay

ローリング、タンブリング、およびスライディング特徴量

負の値ではない必要があります。 datetime.timedelta

start_time

ローリング、タンブリング、スライディング

指定する必要があるのは datetime.datetime

フィールド

サポートされているウィンドウ

制約

delay

ローリング、タンブリング、スライディング

負の値ではない必要があります。 datetime.timedelta

offset

タンブリングおよびスライディング

0以上であり、期間*よりも短い必要があります。

SourceLateness.settling_delay

ローリング、タンブリング、およびスライディング特徴量

負の値ではない必要があります。 datetime.timedelta

start_time

ローリング、タンブリング、スライディング

指定する必要があるのは datetime.datetime

*期間: タンブリングウィンドウの場合、オフセットはwindow_durationより短くする必要があります。スライディングウィンドウの場合、slide_duration より短くする必要があります。

開始時刻​

機能が出力を行える最も早いイベント時間境界をUTCで設定するには、start_timeを使用します。境界は包括的です。start_timeは出力をゲートします。これにより、ウィンドウが読み取る履歴上のソース行が制限されることはなく、ウィンドウの配置が変更されることもありません。start_timeが2つの整列された境界の間にある場合、最初に適格となる固定ウィンドウ出力は次の境界になります。

start_timeを使用すると、ソースで完全なウィンドウ期間が経過する前に、固定期間ウィンドウを出力できます。これらの早期出力では、利用可能なソース履歴が使用されます。たとえば、2024年1月1日にデータが開始するソースに対して、1年間の window_duration と1日間の slide_duration を持つスライディング ウィンドウを考えてみます。

  • start_timeがない場合、完全な1年間のウィンドウを形成できるようになった時点で、この機能は2025年1月1日に初めて出力されます。
  • start_timeを2024年8月21日に設定すると、この機能は2024年8月21日に初めて発行されます。その出力は、2024年1月1日以降の利用可能なソース履歴のみをカバーしています。ウィンドウは2025年1月1日に完全な1年スパンに達し、それ以降は完全な出力を生成します。

start_timeはウィンドウの配置を変更しないため、配置された2つの境界の間の値によって新しい境界が作成されることはありません。UTCの深夜0時に日次境界を持つタンブリングウィンドウの場合、06:00 UTCのstart_timeは、次の深夜0時の境界で最初に出力されます。境界上に正確に収まるstart_timeは、境界が包括的であるため、その境界で出力されます。

start_timeが設定されていない場合、タンブリングウィンドウおよび固定期間のスライディングウィンドウは、完全なウィンドウを形成できるようになった後の整列された境界で最初に発信されます。ライフタイムスライディングウィンドウおよびローリングウィンドウは、適格なソースデータが存在しそうになるとすぐに発信されます。

注記

start_time ローリング、タンブリング、またはスライディングウィンドウでDeltaTableSourceを使用するバッチ特徴量に対してサポートされています。StreamSourceまたはSawtoothWindowではサポートされていません。

例えば:

Python
from datetime import datetime, timedelta
from databricks.feature_engineering.entities import SlidingWindow

window = SlidingWindow(
window_duration=timedelta(days=365),
slide_duration=timedelta(days=1),
start_time=datetime(2024, 8, 21),
)

ローリング ウィンドウ​

注記

RollingWindow 以前はContinuousWindowという名前でした。以前のSDKバージョンから移行している場合は、インポートをそれに応じて更新してください。

ローリング ウィンドウは最新のリアルタイム集計であり、通常はストリーミング データで使用されます。ストリーミングパイプラインでは、イベントの入力または出力など、固定長ウィンドウの内容が変更された場合にのみ、ローリングウィンドウは新しい行を出力します。トレーニングパイプラインでローリングウィンドウ特徴量が使用される場合、特定のイベントのTimestampの直前にある固定長ウィンドウ期間を使用して、ソースデータに対して正確な時点の特徴量計算が実行されます。これは、オンライン/オフラインのずれやデータ漏洩を防ぐのに役立ちます。T時点の特徴量は、[T − 期間, T)からイベントを集約します。

Python
class RollingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
start_time: Optional[datetime.datetime] = None

次の表に、ローリングウィンドウのパラメーターを示します。ウィンドウの開始時刻と終了時刻は、これらのパラメーターに基づいて次のように決まります。

  • 開始時刻:evaluation_time - window_duration - delay(含む)
  • 終了時刻:evaluation_time - delay (排他的)

パラメーター

制約

delay (任意)

0 以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。SourceLateness.settling_delay を使用して、ストリーム内のソース到着遅延の一貫したベースラインをモデル化します。

window_duration

0より大きい必要があります。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

パラメーター

制約

delay (任意)

0 以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。SourceLateness.settling_delay を使用して、ストリーム内のソース到着遅延の一貫したベースラインをモデル化します。

window_duration

0より大きい必要があります。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

Python
from databricks.feature_engineering.entities import RollingWindow
from datetime import timedelta

# Look back 7 days from evaluation time
window = RollingWindow(window_duration=timedelta(days=7))

以下のコードを使用して、遅延を伴うローリングウィンドウを定義します。

Python
# Compute a 7-day value as of one day before the evaluation time
window = RollingWindow(
window_duration=timedelta(days=7),
delay=timedelta(days=1)
)

ローリングウィンドウの例​

  • window_duration=timedelta(days=7)現在の評価時間で終了する7日間のルックバックウィンドウを作成します。7日目の午後2時のイベントの場合、これは0日目の午後2時から7日目の午後2時までの(7日目の午後2時は含まず)すべてのイベントを含みます。

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30):これにより、評価時間の 30 分前に終了する 1 時間のルックバックウィンドウが作成されます。午後3時のイベントの場合、これには午後1時30分から午後2時30分まで(午後2時30分は含まない)のすべてのイベントが含まれます。

Lastを使用して最新値の新しさを制限する​

最新の値が期間限定でのみ有効な場合は、Last と RollingWindow を組み合わせます。評価時には、この間隔内で最新のTimestampを持つ行の値が特徴量によって返されます。

[evaluation_time - delay - window_duration, evaluation_time - delay)

間隔内の最新行に null 値が含まれている場合、特徴量は null を返します。null 入力値を除外する場合は、ソースに filter_condition を設定します。

この組み合わせは ColumnSelection と異なります。ColumnSelection は、経過時間に基づいて有効期限切れにすることなく、観測された最新の NULL ではない値を返します。

バッチ特徴量の場合、この組み合わせには特別なオンライン専用のマテリアライゼーションモードがあります。DeltaTableSource、Last、RollingWindow、およびTableTriggerのみをサポートします。新鮮度で境界が設定された最新値のマテリアライズを参照してください。

タンブリングウィンドウ​

タンブリングウィンドウを使用して定義された特徴量の場合、集計は、スライド間隔で進む所定の固定長ウィンドウで計算され、時間を完全に分割する重複しないウィンドウが生成されます。その結果、ソース内の各イベントは正確に1つのウィンドウに寄与します。時刻tの特徴量は、t (排他的) 以前に終了するウィンドウからデータを集計します。ウィンドウはUnixエポックで開始します。

Python
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
offset: Optional[datetime.timedelta] = None
start_time: Optional[datetime.datetime] = None

次の表に、タンブリングウィンドウのパラメーターを示します。

パラメーター

制約

window_duration

0より大きい必要があります。

delay (任意)

0以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。

offset (任意)

0以上で、window_durationより小さくする必要があります。ウィンドウの境界を UTC の深夜 0 時からシフトします。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

パラメーター

制約

window_duration

0より大きい必要があります。

delay (任意)

0以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。

offset (任意)

0以上で、window_durationより小さくする必要があります。ウィンドウの境界を UTC の深夜 0 時からシフトします。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

Python
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

window = TumblingWindow(
window_duration=timedelta(days=1),
delay=timedelta(hours=2),
offset=timedelta(hours=22),
)

タンブリングウィンドウの例​

  • window_duration=timedelta(days=5)これは、それぞれ5日間の事前に決定された固定長のウィンドウを作成します。例:ウィンドウ #1 は Day 0 から Day 4 にわたり、ウィンドウ #2 は Day 5 から Day 9 にわたり、ウィンドウ #3 は Day 10 から Day 14 にわたり、などです。具体的には、ウィンドウ #1 は、Day 0 の 00:00:00.00 に開始する Timestamp を持つすべてのイベントから、Day 5 の Timestamp 00:00:00.00 を持つイベントまで(ただし、含まず)を含みます。各イベントは厳密に1つのウィンドウに属します。

スライディングウィンドウ​

スライディングウィンドウを使用して定義された特徴量の場合、集計はスライド間隔で進むウィンドウ上で計算されます。スライディングウィンドウには、固定期間またはライフタイム期間のいずれかを設定できます。固定期間ウィンドウは重複するため、各ソースイベントは複数のウィンドウの集計に寄与する可能性があります。ライフタイムウィンドウには、ウィンドウ終了前のすべてのソースイベントが含まれます。時刻tの特徴量は、t(排他的)以前に終了するウィンドウからのデータを集計します。ウィンドウはUnixエポックに準拠しています。

Python
class SlidingWindow(TimeWindow):
window_duration: Optional[datetime.timedelta]
slide_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
offset: Optional[datetime.timedelta] = None
start_time: Optional[datetime.datetime] = None

次の表に、スライディングウィンドウのパラメーターを示します。

パラメーター

制約

window_duration

固定期間ウィンドウの場合は、正の値を設定する必要があります。有効期間ウィンドウにはNoneを設定します。

slide_duration

正の値を設定してください。固定期間ウィンドウの場合、window_durationより短い必要があります。

delay (任意)

0以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。

offset (任意)

0以上で、slide_durationより小さくする必要があります。ウィンドウの境界を UTC の深夜 0 時からシフトします。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

パラメーター

制約

window_duration

固定期間ウィンドウの場合は、正の値を設定する必要があります。有効期間ウィンドウにはNoneを設定します。

slide_duration

正の値を設定してください。固定期間ウィンドウの場合、window_durationより短い必要があります。

delay (任意)

0以上である必要があります。評価Timestampから分析ウィンドウを後方にシフトします。

offset (任意)

0以上で、slide_durationより小さくする必要があります。ウィンドウの境界を UTC の深夜 0 時からシフトします。

start_time (任意)

機能が出力の発行を開始できる最も早いイベント時刻の境界。

Python
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta

window = SlidingWindow(
window_duration=timedelta(days=7),
slide_duration=timedelta(days=1),
delay=timedelta(hours=2),
offset=timedelta(hours=22),
)

スライディングウィンドウの例​

  • window_duration=timedelta(days=5), slide_duration=timedelta(days=1)これにより、毎回1日ずつ進む重複する5日間のウィンドウが作成されます。例: ウィンドウ #1は0日目から4日目まで、ウィンドウ #2は1日目から5日目まで、ウィンドウ #3は2日目から6日目までといった具合です。各ウィンドウには、開始日の00:00:00.00から終了日の00:00:00.00まで (ただしは含まない) のイベントが含まれます。ウィンドウが重複しているため、単一のイベントが複数のウィンドウに属することができます (この例では、各イベントは最大5つの異なるウィンドウに属します)。

有効期間ウィンドウ​

ライフタイムウィンドウを作成するには、window_duration=Noneを設定します。各スライド境界において、特徴量はその境界より前のTimestampを持つエンティティのすべてのソースイベントを集計します。例えば、1日のスライドでは、1日1回累積値が生成されます。

有効期間ウィンドウはSlidingWindowでのみサポートされています。RollingWindowおよびTumblingWindowには、有限のwindow_durationが必要です。

注記

ライフタイムウィンドウには、window_duration=Noneとワークスペースの有効化をサポートするdatabricks-feature-engineeringクライアントバージョンが必要です。以前のクライアントバージョンでは、この構文はサポートされていません。

Python
from datetime import timedelta
from databricks.feature_engineering.entities import (
AggregationFunction,
DeltaTableSource,
Feature,
SlidingWindow,
Sum,
)

lifetime_spend = Feature(
source=DeltaTableSource(
catalog_name="main",
schema_name="store",
table_name="transactions",
),
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(
Sum(input="amount"),
SlidingWindow(
window_duration=None,
slide_duration=timedelta(days=1),
),
),
name="lifetime_spend",
)

ソートゥース ウィンドウ​

備考

ベータ版

SawtoothWindow はベータ版です。

Sawtooth ウィンドウは、最近のイベントに対する非常に新鮮な更新と、履歴データの毎日のコンパクション(圧縮)をサポートする集計です。そのトレーリング(古い)エッジは固定された日単位のステップで進み、リーディング(最近の)エッジは最新のイベントに合わせて更新されるため、有効なウィンドウの長さは1日を通して「のこぎり状」に変化します。ウィンドウの大部分はストリームの取り込みテーブルのデータから提供され、直近の2日間のみがライブストリームから提供されます。これは、新鮮な更新に対する応答性を維持しながら、長時間ウィンドウ(数年単位までスケール可能)を効率的にコンピュートするための妥協案です。

のこぎり型ウィンドウ:リーディングエッジが最新のイベントを追跡し、トレーリングエッジが日単位のステップで進むため、対象となるウィンドウが毎日「のこぎりのように」移動します。

Sawtooth ウィンドウは、バッチとストリーミングのハイブリッドパス上でマテリアライズされます。バッチパイプラインはウィンドウの大部分を保持し、ストリーミングパイプラインは最新のデータをリアルタイムで維持します。これら2つは読み取り時にMergeされるため、モデルやサービングコンシューマーからは単一の特徴量として扱われます。

ウィンドウの履歴部分はバッチパイプラインによってコンピュートされるため、ウィンドウが数か月または数年にわたる場合でも、マテリアライズ開始直後にソートゥース特徴量を提供できるようになります。ローリング ウィンドウは、その全ウィンドウ期間が経過した後にのみ完了します。最小の window_duration は 2 日(強制される下限)より大きければなりません。Databricks では、7 日を超える期間にはソートゥース ウィンドウを推奨しています。2日以上7日までのウィンドウについては、ローリング ウィンドウの固定長の精度と、ソートゥース ウィンドウのより高速な本番運用への対応のどちらかを選択します。

注記

ソートゥース特徴量は、すでに存在する履歴に依存します。ストリームの取り込みテーブルには、少なくともウィンドウの全期間をカバーするデータが含まれている必要があります。含まれていない場合、計算されたウィンドウが不完全になります。2 日が完全に経過する前は、特徴量にはこれまでにマテリアライズされたデータのみが反映されます。完全に 2 日が経過するまでは、本番運用で特徴量を提供することはお勧めしません。空のウィンドウに対する集計では、Sum および Count に対しては 0 が返され、Avg、Min、Max、First、Last、VarPop、VarSamp、StddevPop、および StddevSamp に対しては null が返されます。

Sawtooth 特徴量の準備ができているかどうかを確認するには、カタログエクスプローラーで「Feature View」を開きます。マテリアライズ済み特徴量セクションでは、特徴量の最終マテリアライズ時刻が進み、ステータスが成功を示した時点で、バッチバックフィルは完了です。ストリーミング部分は、Lakeflow 宣言型パイプラインによってマテリアライズされます。Feature View が検証を通過すると、マテリアライズ済み特徴量はそのパイプラインにリンクされ、そこでランステータスを監視できます。

Sawtooth ウィンドウには StreamSource が必要であり、StreamingMode を使用してマテリアライズされます。

Python
class SawtoothWindow(TimeWindow):
window_duration: datetime.timedelta

Sawtooth ウィンドウのエッジは、ローリングウィンドウとは異なる動きをします。リーディングエッジは最新のイベントを追跡し、トレーリングエッジは継続的ではなく1日1回進みます。毎日18:00 UTCの固定カットオフで、トレーリングエッジはその日のUTC深夜の境界まで進みます。その結果、有効なウィンドウは window_duration よりわずかに長くなり、次のカットオフで1日分戻るまで、1日を通して増加します。トレーニングとサービングは同じ18:00 UTCのカットオフを使用するため、オフラインのトレーニングとオンラインのサービングは一貫性を保ちます。

パラメーター

制約

window_duration

2日より長く設定する必要があります。日単位の整数ではない期間(例:timedelta(days=3, minutes=15))も許可されますが、ウィンドウは引き続き日単位の粒度で更新されます。

パラメーター

制約

window_duration

2日より長く設定する必要があります。日単位の整数ではない期間(例:timedelta(days=3, minutes=15))も許可されますが、ウィンドウは引き続き日単位の粒度で更新されます。

ソートゥースウィンドウでは、Sum、Avg、Count、Min、Max、First、Last、VarPop、VarSamp、StddevPop、および StddevSamp の集計関数がサポートされています。

ノコギリ波ウィンドウの例​

次の例は、ユーザーのトランザクションの7日間のカウントを示しています。リーディングエッジは現在のイベントを追跡し、トレーリングエッジは1日ずつ前進します。3月10日のイベントの場合、ウィンドウは約3月3日まで遡ります。3月10日が経過するにつれて、リーディングエッジは前進し続けますが、トレーリングエッジは保持されるため、カバーされる期間は拡大します。その後、3月11日の開始時に、トレーリングエッジは約3月4日まで進みます。有効なウィンドウは常に7日間より少し長くなります。直近の2日間はライブストリームから提供され、それ以前の日付はストリームの取り込みテーブルから提供されます。

Python
from databricks.feature_engineering.entities import SawtoothWindow
from datetime import timedelta

# 7-day window kept continuously fresh with streaming data
window = SawtoothWindow(window_duration=timedelta(days=7))

ノコギリ波ウィンドウの制限​

  • delay パラメーターはサポートされていません。
  • SourceLateness.settling_delay はサポートされていません。
  • Sum、Avg、Count、Min、Max、First、Last、VarPop、VarSamp、StddevPop、および StddevSamp 以外の集計関数はサポートされていません(例:ApproxCountDistinct、ApproxPercentile、FirstN、LastN、FirstDistinct、LastDistinct)。
  • Sawtooth ウィンドウには StreamSource が必要です。DeltaTableSource はサポートされていません。

Materialization Trigger​

トリガーは、実体化パイプラインが実行されるタイミングを制御します。Triggerタイプはフィーチャータイプによって異なります。

CronSchedule​

バッチ集計機能には CronSchedule を使用します。By default, Databricks は集計ウィンドウからスケジュールを派生させます。派生スケジュールでは、ウィンドウ期間、ウィンドウdelayおよびoffset、ならびにソースsettling_delayが考慮されるため、ソースデータの完了が予想される前にランがウィンドウをパブリッシュすることはありません。派生スケジュールは、タンブリングウィンドウおよびスライディングウィンドウをサポートしています。

派生スケジュールをリクエストするには、cron 式を省略します。CronSchedule() と明示的な CronSchedule(mode=CronScheduleMode.DERIVED) 形式は同等です:

Python
from databricks.feature_engineering.entities import (
CronSchedule,
CronScheduleMode,
)

trigger = CronSchedule(mode=CronScheduleMode.DERIVED)

quartz_cron_expression に CronScheduleMode.DERIVED を設定しないでください。マテリアライズされた特徴量を取得すると、返されるスケジュールに Databricks が計算した cron 式が含まれる場合があります。

スケジュールを直接制御するには、Quartz cron式を指定します。式を指定すると、CronScheduleMode.MANUAL が推論されます。

Python
from databricks.feature_engineering.entities import CronSchedule

trigger = CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
)

TableTrigger​

ColumnSelection 特徴量または DeltaTableSource に裏打ちされた集約特徴量 (AggregationFunction) には、TableTrigger を使用します。アップストリームの Delta テーブルが新しい commit を受け取るたびに、パイプラインが実行されます。

集計機能の場合、パイプラインはすべての commit ごとに実行されないようにスロットルされます。パイプラインは、機能のウィンドウ長の半分につき最大 1 回実行されますが、5 分より短い間隔で実行されることはありません。例えば、1 時間のタンブリングウィンドウを持つ機能は 30 分ごとに最大 1 回実行され、8 時間のウィンドウを持つ機能は 4 時間ごとに最大 1 回実行されます。ウィンドウの半分が 5 分未満の場合、5 分という下限が適用されるため、10 分以下のウィンドウは 5 分ごとに最大 1 回実行されます。ウィンドウが 5 分未満の集計機能では TableTrigger を使用できません。代わりにストリーミング Trigger を使用してください。

Python
from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode​

StreamSourceによってサポートされる特徴量にはStreamingModeを使用します。パイプラインは、継続的なストリーミングパイプラインとして実行されます。

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

fe = FeatureEngineeringClient()

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

streaming_feature = fe.create_feature(
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)),
),
catalog_name="my_catalog",
schema_name="my_schema",
name="user_purchase_sum",
)

fe.materialize_features(
features=[streaming_feature],
online_config=OnlineStoreConfig(
catalog_name="my_catalog",
schema_name="my_schema",
table_name_prefix="streaming_features_serving",
online_store_name="feature_store_online",
),
trigger=StreamingMode(),
)

Triggerの選択​

各特徴量は1つの Trigger を使用します。特徴量タイプ別のオプションは次のとおりです:

特徴量タイプ

トリガー

実行するとき

集計(AggregationFunction)元 DeltaTableSource

CronSchedule

派生スケジュールまたは手動スケジュール

集計(AggregationFunction)元 DeltaTableSource

TableTrigger

各ソーステーブルcommit時

ColumnSelection (DeltaTableSource から)

TableTrigger

各ソーステーブルcommit時

からの特徴量 StreamSource

StreamingMode

連続ストリーミング

特徴量タイプ

トリガー

実行するとき

集計(AggregationFunction)元 DeltaTableSource

CronSchedule

派生スケジュールまたは手動スケジュール

集計(AggregationFunction)元 DeltaTableSource

TableTrigger

各ソーステーブルcommit時

ColumnSelection (DeltaTableSource から)

TableTrigger

各ソーステーブルcommit時

からの特徴量 StreamSource

StreamingMode

連続ストリーミング

単一のmaterialize_features呼び出しで、異なるTriggerタイプを必要とする特徴量をマテリアライズすることはできません。代わりに、個別の呼び出しを発行します。